Merge remote-tracking branch 'origin/master' into master-jb

This commit is contained in:
hzj
2026-07-22 11:02:27 +08:00
35 changed files with 967 additions and 329 deletions

View File

@@ -16,6 +16,7 @@ import com.njcn.common.pojo.annotation.OperateInfo;
import com.njcn.common.pojo.enums.common.LogEnum;
import com.njcn.common.pojo.exception.BusinessException;
import com.njcn.csdevice.api.CsLineFeignClient;
import com.njcn.csdevice.pojo.dto.CsLineDTO;
import com.njcn.device.biz.commApi.CommTerminalGeneralClient;
import com.njcn.device.biz.pojo.dto.DeptGetChildrenMoreDTO;
import com.njcn.device.biz.pojo.dto.DeptGetDeviceDTO;
@@ -23,6 +24,7 @@ import com.njcn.device.biz.pojo.dto.DeptGetSubStationDTO;
import com.njcn.device.biz.pojo.dto.LineDevGetDTO;
import com.njcn.device.biz.pojo.param.DeptGetLineParam;
import com.njcn.device.pq.api.DeptLineFeignClient;
import com.njcn.redis.utils.RedisUtil;
import com.njcn.user.api.DeptFeignClient;
import com.njcn.user.pojo.po.Dept;
import com.njcn.web.controller.BaseController;
@@ -73,6 +75,8 @@ public class ExecutionCenter extends BaseController {
private CsLineFeignClient csLineFeignClient;
@Resource
private FlowAsyncService flowService;
@Resource
private RedisUtil redisUtil;
/***
* 1、校验非全链执行时tagNames节点标签集合必须为非空否则提示---无可执行节点
@@ -143,9 +147,27 @@ public class ExecutionCenter extends BaseController {
String methodDescribe = getMethodDescribe("wlMeasurementPointExecutor");
//手动判断参数是否合法,
CalculatedParam calculatedParam = judgeExecuteParam(baseParam);
//缓存监测点、接线方式数据
List<CsLineDTO> csLineDTOList = csLineFeignClient.getAllLineDetail().getData();
List<String> lineIdList = new ArrayList<>();
Map<String, Integer> lineConTypeMap = new HashMap<>();
if (CollectionUtils.isNotEmpty(csLineDTOList)) {
lineIdList = csLineDTOList.stream().map(CsLineDTO::getLineId).collect(Collectors.toList());
//接线方式(0-星型 1-角型 2-V型)
lineConTypeMap = csLineDTOList.stream()
.filter(dto -> dto.getLineId() != null)
.collect(Collectors.toMap(
CsLineDTO::getLineId,
dto -> dto.getConType() != null ? dto.getConType() : 0,
(v1, v2) -> v1
));
redisUtil.saveByKeyWithExpire("wlLineDetail", lineConTypeMap, 3600L*3);
}
// 测点索引
if (CollectionUtils.isEmpty(calculatedParam.getIdList())) {
calculatedParam.setIdList(csLineFeignClient.getAllLine().getData());
calculatedParam.setIdList(lineIdList);
} else {
calculatedParam.setIdList(calculatedParam.getIdList());
}
LiteflowResponse liteflowResponse;
if (baseParam.isRepair()) {

View File

@@ -2,7 +2,6 @@ package com.njcn.algorithm.service.line;
import com.njcn.algorithm.pojo.bo.CalculatedParam;
import java.util.concurrent.CompletableFuture;
/**
* @author xy

View File

@@ -52,8 +52,8 @@ public class DataCleanServiceImpl implements IDataCleanService {
private static final Logger logger = LoggerFactory.getLogger(DataCleanServiceImpl.class);
@Value("${line.num}")
private Integer NUM = 100;
@Value("${line.num:10}")
private Integer NUM;
@Resource
private DataVFeignClient dataVFeignClient;
@@ -101,7 +101,7 @@ public class DataCleanServiceImpl implements IDataCleanService {
public void dataQualityCleanHandler(CalculatedParam calculatedParam) {
MemorySizeUtil.getNowMemory();
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS");
logger.info("{},数据质量清洗算法执行=====》", LocalDateTime.now());
logger.info("{},{}数据质量清洗算法执行=====》", LocalDateTime.now(),calculatedParam.getDataDate());
//获取监测点的统计间隔
List<String> listOfString = (List<String>) (List<?>) calculatedParam.getIdList();
List<LineDetailDataVO> lineDetailDataVOS = lineFeignClient.getLineDetailList(listOfString).getData();
@@ -290,7 +290,7 @@ public class DataCleanServiceImpl implements IDataCleanService {
DictData dip = dicDataFeignClient.getDicDataByCodeAndType(DicDataEnum.VOLTAGE_DIP.getCode(), DicDataTypeEnum.EVENT_STATIS.getCode()).getData();
DictData rise = dicDataFeignClient.getDicDataByCodeAndType(DicDataEnum.VOLTAGE_RISE.getCode(), DicDataTypeEnum.EVENT_STATIS.getCode()).getData();
MemorySizeUtil.getNowMemory();
logger.info("{},原始表数据清洗=====》", LocalDateTime.now());
logger.info("{},{}原始表数据清洗=====》", LocalDateTime.now(),calculatedParam.getDataDate());
//获取标准
Map<String, List<PqReasonableRangeDto>> map = getStandardData();
//获取监测点台账信息
@@ -304,7 +304,7 @@ public class DataCleanServiceImpl implements IDataCleanService {
lineParam.setStartTime(TimeUtils.getBeginOfDay(calculatedParam.getDataDate()));
lineParam.setEndTime(TimeUtils.getEndOfDay(calculatedParam.getDataDate()));
for (int i = 0; i < lineDetail.size(); i++) {
logger.info( calculatedParam.getDataDate()+"总数据:" + lineDetail.size() + "=====》当前第" + (i + 1));
// logger.info( calculatedParam.getDataDate()+"总数据:" + lineDetail.size() + "=====》当前第" + (i + 1));
LineDetailVO.Detail item = lineDetail.get(i);
flowService.lineDataClean(item, map, calculatedParam.getDataDate(), dip, rise,lineDetail.size(),(i + 1));
}

View File

@@ -3,6 +3,7 @@ package com.njcn.algorithm.serviceimpl.line;
import cn.hutool.core.collection.CollUtil;
import com.njcn.algorithm.pojo.bo.CalculatedParam;
import com.njcn.algorithm.service.line.IDataComAssService;
import com.njcn.algorithm.utils.MemorySizeUtil;
import com.njcn.dataProcess.api.DataComAssFeignClient;
import com.njcn.dataProcess.api.DataFlickerFeignClient;
import com.njcn.dataProcess.api.DataVFeignClient;
@@ -38,8 +39,8 @@ import java.util.stream.Collectors;
public class DataComAssServiceImpl implements IDataComAssService {
private static final Logger logger = LoggerFactory.getLogger(DayDataServiceImpl.class);
@Value("${line.num}")
private Integer NUM = 100;
@Resource
private DataVFeignClient dataVFeignClient;
@@ -50,10 +51,17 @@ public class DataComAssServiceImpl implements IDataComAssService {
@Resource
private PqDataVerifyFeignClient pqDataVerifyFeignClient;
/**
* 查询配置 辽宁版本比较特殊没有CP95值采用平均值进行判断
*/
@Value("${version.used:master}")
private String versionUsed;
@Override
public void dataComAssHandler(CalculatedParam calculatedParam) {
List<DataComassesDPO> list = new ArrayList<>();
logger.info("{},r_stat_comasses_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}r_stat_comasses_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
lineParam.setStartTime(TimeUtils.getBeginOfDay(calculatedParam.getDataDate()));
lineParam.setEndTime(TimeUtils.getEndOfDay(calculatedParam.getDataDate()));
@@ -136,7 +144,7 @@ public class DataComAssServiceImpl implements IDataComAssService {
BigDecimal hundred = BigDecimal.valueOf(100);
//************************************电压偏差********************************************
lineParam.setPhasicType(Arrays.asList("A", "B", "C"));
lineParam.setValueType(Arrays.asList("AVG"));
lineParam.setValueType(Collections.singletonList("AVG"));
lineParam.setColumnName("vu_dev");
lineParam.setGe("10");
Integer vuDev1 = dataVFeignClient.getColumnNameCountRawData(lineParam).getData();
@@ -167,6 +175,7 @@ public class DataComAssServiceImpl implements IDataComAssService {
//************************************频率偏差********************************************
lineParam.setColumnName("freq_dev");
lineParam.setGe("0.3");
lineParam.setLt(null);
Integer freqDev1 = dataVFeignClient.getColumnNameCountRawData(lineParam).getData();
lineParam.setGe("0.2");
@@ -194,8 +203,15 @@ public class DataComAssServiceImpl implements IDataComAssService {
}
//************************************谐波畸变率********************************************
lineParam.setColumnName("v_thd");
lineParam.setValueType(Arrays.asList("CP95"));
if("liaoning".equals(versionUsed)){
lineParam.setValueType(Collections.singletonList("AVG"));
}else {
lineParam.setValueType(Collections.singletonList("CP95"));
}
lineParam.setGe("6");
lineParam.setLt(null);
Integer vThd1 = dataVFeignClient.getColumnNameCountRawData(lineParam).getData();
lineParam.setGe("4");
@@ -224,7 +240,12 @@ public class DataComAssServiceImpl implements IDataComAssService {
//************************************三相电压不平衡度********************************************
lineParam.setColumnName("v_unbalance");
if("liaoning".equals(versionUsed)){
lineParam.setValueType(Collections.singletonList("MAX"));
}
lineParam.setPhasicType(Collections.singletonList("T"));
lineParam.setGe("4");
lineParam.setLt(null);
Integer vUnbalance1 = dataVFeignClient.getColumnNameCountRawData(lineParam).getData();
lineParam.setGe("2");
@@ -253,7 +274,9 @@ public class DataComAssServiceImpl implements IDataComAssService {
//************************************电压波动(短时闪变)********************************************
lineParam.setColumnName("pst");
lineParam.setValueType(null);
lineParam.setPhasicType(Arrays.asList("A", "B", "C"));
lineParam.setGe("0.8");
lineParam.setLt("1");
Integer plt1 = dataFlickerFeignClient.getColumnNameCountRawData(lineParam).getData();
lineParam.setGe("0.6");
@@ -273,11 +296,11 @@ public class DataComAssServiceImpl implements IDataComAssService {
Integer plt5 = dataFlickerFeignClient.getColumnNameCountRawData(lineParam).getData();
BigDecimal pltAll = BigDecimal.valueOf(plt1 + plt2 + plt3 + plt4 + plt5);
if (pltAll.compareTo(BigDecimal.ZERO) != 0) {
outMap.put("data_plt1", BigDecimal.valueOf(plt1).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_plt2", BigDecimal.valueOf(plt2).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_plt3", BigDecimal.valueOf(plt3).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_plt4", BigDecimal.valueOf(plt4).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_plt5", BigDecimal.valueOf(plt5).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_pst1", BigDecimal.valueOf(plt1).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_pst2", BigDecimal.valueOf(plt2).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_pst3", BigDecimal.valueOf(plt3).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_pst4", BigDecimal.valueOf(plt4).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
outMap.put("data_pst5", BigDecimal.valueOf(plt5).multiply(hundred).divide(pltAll, 3, RoundingMode.HALF_UP));
}
if (!CollUtil.isEmpty(outMap)) {

View File

@@ -35,8 +35,9 @@ public class DayDataServiceImpl implements IDayDataService {
private static final Logger logger = LoggerFactory.getLogger(DayDataServiceImpl.class);
@Value("${line.num}")
private Integer NUM = 100;
@Value("${line.num:10}")
private Integer NUM;
@Resource
private DataVFeignClient dataVFeignClient;
@Resource
@@ -75,7 +76,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataVHandler(CalculatedParam calculatedParam) {
logger.info("{},dataV表转r_stat_data_v_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataV表转r_stat_data_v_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataVDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -91,7 +93,6 @@ public class DayDataServiceImpl implements IDayDataService {
//获取原始数据
List<CommonMinuteDto> partList = dataVFeignClient.getBaseData(lineParam).getData();
if (CollUtil.isNotEmpty(partList)) {
logger.info("{}dataV集合大小为>>>>>>>>>>>>{}",lineParam.getStartTime(),MemorySizeUtil.getObjectSize(partList));
partList.forEach(item->{
//相别
List<CommonMinuteDto.PhasicType> phasicTypeList = item.getPhasicTypeList();
@@ -125,7 +126,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataIHandler(CalculatedParam calculatedParam) {
logger.info("{},dataI表转r_stat_data_i_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataI表转r_stat_data_i_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataIDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -173,7 +175,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataFlickerHandler(CalculatedParam calculatedParam) {
logger.info("{},dataFlicker表转r_stat_data_flicker_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataFlicker表转r_stat_data_flicker_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<String> valueList = Arrays.asList("AVG","MAX","MIN","CP95");
List<DataFlickerDto> result = new ArrayList<>();
//远程接口获取分钟数据
@@ -222,7 +225,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataFlucHandler(CalculatedParam calculatedParam) {
logger.info("{},dataFluc表转r_stat_data_fluc_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataFluc表转r_stat_data_fluc_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<String> valueList = Arrays.asList("AVG","MAX","MIN","CP95");
List<DataFlucDto> result = new ArrayList<>();
//远程接口获取分钟数据
@@ -271,7 +275,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataHarmPhasicIHandler(CalculatedParam calculatedParam) {
logger.info("{},dataHarmPhasicI表转r_stat_data_harmphasic_i_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataHarmPhasicI表转r_stat_data_harmphasic_i_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataHarmPhasicIDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -319,7 +324,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataHarmPhasicVHandler(CalculatedParam calculatedParam) {
logger.info("{},dataHarmPhasicV表转r_stat_data_harmphasic_v_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataHarmPhasicV表转r_stat_data_harmphasic_v_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataHarmPhasicVDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -367,7 +373,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataHarmPowerPHandler(CalculatedParam calculatedParam) {
logger.info("{},dataHarmPowerP表转r_stat_data_harmpower_p_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataHarmPowerP表转r_stat_data_harmpower_p_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataHarmPowerPDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -415,7 +422,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataHarmPowerQHandler(CalculatedParam calculatedParam) {
logger.info("{},dataHarmPowerQ表转r_stat_data_harmpower_q_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataHarmPowerQ表转r_stat_data_harmpower_q_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataHarmPowerQDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -463,7 +471,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataHarmPowerSHandler(CalculatedParam calculatedParam) {
logger.info("{},dataHarmPowerS表转r_stat_data_harmpower_s_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataHarmPowerS表转r_stat_data_harmpower_s_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataHarmPowerSDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -511,7 +520,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataHarmRateIHandler(CalculatedParam calculatedParam) {
logger.info("{},dataHarmRateI表转r_stat_data_harmRate_i_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataHarmRateI表转r_stat_data_harmRate_i_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataHarmRateIDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -559,7 +569,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataHarmRateVHandler(CalculatedParam calculatedParam) {
logger.info("{},dataHarmRateV表转r_stat_data_harmRate_v_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataHarmRateV表转r_stat_data_harmRate_v_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataHarmRateVDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -607,7 +618,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataInHarmIHandler(CalculatedParam calculatedParam) {
logger.info("{},dataInHarmI表转r_stat_data_inharm_i_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataInHarmI表转r_stat_data_inharm_i_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataInHarmIDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -655,7 +667,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataInHarmVHandler(CalculatedParam calculatedParam) {
logger.info("{},dataInHarmV表转r_stat_data_inharm_v_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataInHarmV表转r_stat_data_inharm_v_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataInHarmVDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -703,7 +716,8 @@ public class DayDataServiceImpl implements IDayDataService {
@Override
public void dataPltHandler(CalculatedParam calculatedParam) {
logger.info("{},dataPlt表转r_stat_data_plt_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}dataPlt表转r_stat_data_plt_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataPltDto> result = new ArrayList<>();
List<String> valueList = Arrays.asList("AVG","MAX","MIN","CP95");
//远程接口获取分钟数据

View File

@@ -390,16 +390,16 @@ public class FlowAsyncServiceImpl implements FlowAsyncService {
bak.setState(1);
//存储文件
InputStream reportStream = IoUtil.toStream(new Gson().toJson(resultData), CharsetUtil.UTF_8);
String[] date = dataDate.split("-");
String fileName = fileStorageUtil.uploadStreamSpecifyName(
reportStream
, OssPath.DATA_CLEAN + dataDate + "/"
, OssPath.DATA_CLEAN + date[0] + "/"+ date[1] + "/"+ date[2] + "/"
, item.getLineId() + ".txt");
//存储数据
bak.setFilePath(fileName);
}
pqDataVerifyNewFeignClient.insertData(bak);
resultData=null;
logger.info( dataDate+"总数据:" + size + "=====》当前第" + i+"已完成!");
System.gc();
}

View File

@@ -5,6 +5,7 @@ import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.collection.ListUtil;
import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.LocalDateTimeUtil;
import cn.hutool.core.util.ArrayUtil;
import cn.hutool.core.util.ObjectUtil;
import com.njcn.algorithm.pojo.bo.CalculatedParam;
import com.njcn.algorithm.service.line.IDataCrossingService;
@@ -42,10 +43,7 @@ import java.lang.reflect.Method;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.function.Function;
import java.util.stream.Collectors;
@@ -61,7 +59,7 @@ import java.util.stream.Collectors;
public class IDataCrossingServiceImpl implements IDataCrossingService {
private static final Logger logger = LoggerFactory.getLogger(DayDataServiceImpl.class);
@Value("${line.num}")
@Value("${line.num:10}")
private Integer NUM;
@Resource
@@ -106,29 +104,28 @@ public class IDataCrossingServiceImpl implements IDataCrossingService {
}
Map<String, Overlimit> overLimitMap = overLimitList.stream().collect(Collectors.toMap(Overlimit::getId, Function.identity()));
//以100个监测点分片处理
List<List<String>> pendingIds = ListUtils.partition(lineIds, 1);
ArrayList<String> phase = ListUtil.toList(PhaseType.PHASE_A, PhaseType.PHASE_B, PhaseType.PHASE_C);
MemorySizeUtil.getNowMemory();
List<List<String>> pendingIds = ListUtils.partition(lineIds, 20);
ArrayList<String> phase = ListUtil.toList(PhaseType.PHASE_A, PhaseType.PHASE_B, PhaseType.PHASE_C);
List<CompletableFuture<Void>> futures = new ArrayList<>();
for (int i = 0; i < pendingIds.size(); i++) {
logger.info(calculatedParam.getDataDate()+" 总分区数据:" + pendingIds.size() + "=====》当前第"+(i + 1)+"小分区");
List<String> list = pendingIds.get(i);
// 获取Future
CompletableFuture<Void> future = dataLimitRateAsync.lineDataRate(
calculatedParam.getDataDate(),
list,
phase,
overLimitMap,
pendingIds.size(),
(i + 1),
lineParam.getType()
);
futures.add(future);
for (List<String> pendingId : pendingIds) {
List<CompletableFuture<Void>> futures = new ArrayList<>();
for (String s : pendingId) {
CompletableFuture<Void> future = dataLimitRateAsync.lineDataRate(
calculatedParam.getDataDate(),
Arrays.asList(s),
phase,
overLimitMap,
pendingId.size(),
pendingId.size(),
lineParam.getType()
);
futures.add(future);
}
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
}
// 等待所有任务完成
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
System.err.println("limitRate表转r_stat_limit_rate_d算法开始执行日期为{},执行完成=====》"+calculatedParam.getDataDate());
System.gc();
}
@@ -169,7 +166,8 @@ public class IDataCrossingServiceImpl implements IDataCrossingService {
@Override
public void limitQualifiedDayHandler(CalculatedParam calculatedParam) {
logger.info("{},r_stat_limit_qualified_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}r_stat_limit_qualified_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
lineParam.setStartTime(TimeUtils.getBeginOfDay(calculatedParam.getDataDate()));
lineParam.setEndTime(TimeUtils.getEndOfDay(calculatedParam.getDataDate()));

View File

@@ -5,6 +5,7 @@ import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.LocalDateTimeUtil;
import com.njcn.algorithm.pojo.bo.CalculatedParam;
import com.njcn.algorithm.service.line.IDataIntegrityService;
import com.njcn.algorithm.utils.MemorySizeUtil;
import com.njcn.dataProcess.api.DataIntegrityFeignClient;
import com.njcn.dataProcess.api.DataVFeignClient;
import com.njcn.dataProcess.dto.MeasurementCountDTO;
@@ -42,8 +43,8 @@ public class IDataIntegrityServiceImpl implements IDataIntegrityService {
private static final Logger logger = LoggerFactory.getLogger(IDataIntegrityServiceImpl.class);
@Value("${line.num}")
private Integer NUM = 100;
@Value("${line.num:10}")
private Integer NUM;
@Resource
private CommTerminalGeneralClient commTerminalGeneralClient;
@Resource
@@ -54,7 +55,8 @@ public class IDataIntegrityServiceImpl implements IDataIntegrityService {
@Override
public void dataIntegrity(CalculatedParam<String> calculatedParam) {
logger.info("{},integrity表转r_stat_integrity_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}integrity表转r_stat_integrity_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataIntegrityDto> poList = new ArrayList<>();
List<String> lineIds = calculatedParam.getIdList();
String beginDay = LocalDateTimeUtil.format(

View File

@@ -294,11 +294,11 @@ public class IDataLimitRateAsyncImpl implements IDataLimitRateAsync {
dataLimitRateDetailFeignClient.batchInsertion(detail);
}
}
logger.info(dataDate + " 总分区数据:" + size + "=====》当前第" + i + "小分区已完成!");
// logger.info(dataDate + " 总分区数据:" + size + "=====》当前第" + i + "小分区已完成!");
result = null;
if(i==size){
MemorySizeUtil.getNowMemory();
}
// if(i==size){
// MemorySizeUtil.getNowMemory();
// }
System.gc();
return CompletableFuture.completedFuture(null);

View File

@@ -9,6 +9,7 @@ import cn.hutool.core.util.NumberUtil;
import cn.hutool.core.util.ObjectUtil;
import com.njcn.algorithm.pojo.bo.CalculatedParam;
import com.njcn.algorithm.service.line.IDataOnlineRateService;
import com.njcn.algorithm.utils.MemorySizeUtil;
import com.njcn.dataProcess.api.DataIntegrityFeignClient;
import com.njcn.dataProcess.api.DataOnlineRateFeignClient;
import com.njcn.dataProcess.api.DataVFeignClient;
@@ -48,8 +49,8 @@ public class IDataOnlineRateServiceImpl implements IDataOnlineRateService {
private static final Logger logger = LoggerFactory.getLogger(IDataOnlineRateServiceImpl.class);
@Value("${line.num}")
private Integer NUM = 100;
@Value("${line.num:10}")
private Integer NUM;
private final Integer online = 1;
private final Integer offline = 0;
@@ -69,7 +70,8 @@ public class IDataOnlineRateServiceImpl implements IDataOnlineRateService {
@Override
public void dataOnlineRate(CalculatedParam calculatedParam) {
logger.info("{},onlineRate表转r_stat_onlinerate_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}onlineRate表转r_stat_onlinerate_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
lineParam.setStartTime(TimeUtils.getBeginOfDay(calculatedParam.getDataDate()));
lineParam.setEndTime(TimeUtils.getEndOfDay(calculatedParam.getDataDate()));

View File

@@ -7,6 +7,7 @@ import cn.hutool.core.date.LocalDateTimeUtil;
import cn.hutool.core.util.ObjectUtil;
import com.njcn.algorithm.pojo.bo.CalculatedParam;
import com.njcn.algorithm.service.line.IEventDetailService;
import com.njcn.algorithm.utils.MemorySizeUtil;
import com.njcn.dataProcess.api.EventDetailFeignClient;
import com.njcn.dataProcess.api.RmpEventDetailFeignClient;
import com.njcn.dataProcess.dto.RmpEventDetailDTO;
@@ -44,8 +45,8 @@ import java.util.stream.Collectors;
@RequiredArgsConstructor
public class IEventDetailServiceImpl implements IEventDetailService {
private static final Logger logger = LoggerFactory.getLogger(IDataOnlineRateServiceImpl.class);
@Value("${line.num}")
private Integer NUM = 10;
@Value("${line.num:10}")
private Integer NUM;
@Resource
private RmpEventDetailFeignClient eventDetailFeignClient;
@Resource
@@ -55,7 +56,8 @@ public class IEventDetailServiceImpl implements IEventDetailService {
@Override
public void dataDayHandle(CalculatedParam<String> calculatedParam) {
logger.info("{},r_mp_event_detail_d算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}r_mp_event_detail_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataEventDetailDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -106,7 +108,8 @@ public class IEventDetailServiceImpl implements IEventDetailService {
@Override
public void dataMonthHandle(CalculatedParam<String> calculatedParam) {
logger.info("{},r_mp_event_detail_m算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}r_mp_event_detail_m算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataEventDetailDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -139,7 +142,8 @@ public class IEventDetailServiceImpl implements IEventDetailService {
@Override
public void dataQuarterHandle(CalculatedParam<String> calculatedParam) {
logger.info("{},r_mp_event_detail_q算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}r_mp_event_detail_q算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataEventDetailDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();
@@ -172,7 +176,8 @@ public class IEventDetailServiceImpl implements IEventDetailService {
@Override
public void dataYearHandle(CalculatedParam<String> calculatedParam) {
logger.info("{},r_mp_event_detail_y算法开始=====》", LocalDateTime.now());
MemorySizeUtil.getNowMemory();
logger.info("{},{}r_mp_event_detail_y算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate());
List<DataEventDetailDto> result = new ArrayList<>();
//远程接口获取分钟数据
LineCountEvaluateParam lineParam = new LineCountEvaluateParam();

View File

@@ -49,8 +49,8 @@ import java.util.stream.Stream;
@RequiredArgsConstructor
public class PollutionServiceImpl implements IPollutionService {
@Value("${line.num}")
private Integer NUM = 100;
@Value("${line.num:10}")
private Integer NUM;
@Resource
private PqDataVerifyFeignClient pqDataVerifyFeignClient;

View File

@@ -21,6 +21,7 @@ import com.njcn.device.biz.commApi.CommLineClient;
import com.njcn.device.biz.pojo.dto.LineDTO;
import com.njcn.device.pms.pojo.param.MonitorTerminalParam;
import com.njcn.device.pq.api.LineFeignClient;
import com.njcn.device.pq.api.UserLedgerFeignClient;
import com.njcn.device.pq.pojo.vo.LineDetailDataVO;
import com.njcn.event.api.EventDetailFeignClient;
import com.njcn.event.api.TransientFeignClient;
@@ -65,8 +66,8 @@ import java.util.stream.Collectors;
@RequiredArgsConstructor
public class SpecialAnalysisServiceImpl implements ISpecialAnalysisService {
private Integer NUM=10;
@Value("${line.num:10}")
private Integer NUM;
@Resource
private EventDetailFeignClient eventDetailFeignClient;
@Resource
@@ -84,7 +85,7 @@ public class SpecialAnalysisServiceImpl implements ISpecialAnalysisService {
@Resource
private CommLineClient commLineClient;
@Resource
private UserLedgerOldFeignClient userLedgerFeignClient;
private UserLedgerFeignClient userLedgerFeignClient;
@Resource
private DataLimitRateDetailFeignClient dataLimitRateDetailFeignClient;

View File

@@ -2,10 +2,7 @@ package com.njcn.dataProcess.api;
import com.njcn.common.pojo.constant.ServerInfo;
import com.njcn.common.pojo.response.HttpResult;
import com.njcn.dataProcess.api.fallback.RmpEventFeignClientFallbackFactory;
import com.njcn.dataProcess.api.fallback.SpThroughFeignClientFallbackFactory;
import com.njcn.dataProcess.dto.RmpEventDetailDTO;
import com.njcn.dataProcess.param.LineCountEvaluateParam;
import com.njcn.dataProcess.pojo.dto.RActivePowerRangeDto;
import com.njcn.dataProcess.pojo.dto.SpThroughDto;
import org.springframework.cloud.openfeign.FeignClient;

View File

@@ -214,111 +214,111 @@ public class DataHarmpowerP {
private Double p50;
//在线监测添加的字段
@Column(name = "totPf")
private Double totPf;
@Column(name = "totDf")
private Double totDf;
@Column(name = "totP")
@Column(name = "tot_p")
private Double totP;
@Column(name = "totP1")
@Column(name = "tot_df")
private Double totDf;
@Column(name = "tot_pf")
private Double totPf;
@Column(name = "tot_tp_1")
private Double totP1;
@Column(name = "totP2")
@Column(name = "tot_tp_2")
private Double totP2;
@Column(name = "totP3")
@Column(name = "tot_tp_3")
private Double totP3;
@Column(name = "totP4")
@Column(name = "tot_tp_4")
private Double totP4;
@Column(name = "totP5")
@Column(name = "tot_tp_5")
private Double totP5;
@Column(name = "totP6")
@Column(name = "tot_tp_6")
private Double totP6;
@Column(name = "totP7")
@Column(name = "tot_tp_7")
private Double totP7;
@Column(name = "totP8")
@Column(name = "tot_tp_8")
private Double totP8;
@Column(name = "totP9")
@Column(name = "tot_tp_9")
private Double totP9;
@Column(name = "totP10")
@Column(name = "tot_tp_10")
private Double totP10;
@Column(name = "totP11")
@Column(name = "tot_tp_11")
private Double totP11;
@Column(name = "totP12")
@Column(name = "tot_tp_12")
private Double totP12;
@Column(name = "totP13")
@Column(name = "tot_tp_13")
private Double totP13;
@Column(name = "totP14")
@Column(name = "tot_tp_14")
private Double totP14;
@Column(name = "totP15")
@Column(name = "tot_tp_15")
private Double totP15;
@Column(name = "totP16")
@Column(name = "tot_tp_16")
private Double totP16;
@Column(name = "totP17")
@Column(name = "tot_tp_17")
private Double totP17;
@Column(name = "totP18")
@Column(name = "tot_tp_18")
private Double totP18;
@Column(name = "totP19")
@Column(name = "tot_tp_19")
private Double totP19;
@Column(name = "totP20")
@Column(name = "tot_tp_20")
private Double totP20;
@Column(name = "totP21")
@Column(name = "tot_tp_21")
private Double totP21;
@Column(name = "totP22")
@Column(name = "tot_tp_22")
private Double totP22;
@Column(name = "totP23")
@Column(name = "tot_tp_23")
private Double totP23;
@Column(name = "totP24")
@Column(name = "tot_tp_24")
private Double totP24;
@Column(name = "totP25")
@Column(name = "tot_tp_25")
private Double totP25;
@Column(name = "totP26")
@Column(name = "tot_tp_26")
private Double totP26;
@Column(name = "totP27")
@Column(name = "tot_tp_27")
private Double totP27;
@Column(name = "totP28")
@Column(name = "tot_tp_28")
private Double totP28;
@Column(name = "totP29")
@Column(name = "tot_tp_29")
private Double totP29;
@Column(name = "totP30")
@Column(name = "tot_tp_30")
private Double totP30;
@Column(name = "totP31")
@Column(name = "tot_tp_31")
private Double totP31;
@Column(name = "totP32")
@Column(name = "tot_tp_32")
private Double totP32;
@Column(name = "totP33")
@Column(name = "tot_tp_33")
private Double totP33;
@Column(name = "totP34")
@Column(name = "tot_tp_34")
private Double totP34;
@Column(name = "totP35")
@Column(name = "tot_tp_35")
private Double totP35;
@Column(name = "totP36")
@Column(name = "tot_tp_36")
private Double totP36;
@Column(name = "totP37")
@Column(name = "tot_tp_37")
private Double totP37;
@Column(name = "totP38")
@Column(name = "tot_tp_38")
private Double totP38;
@Column(name = "totP39")
@Column(name = "tot_tp_39")
private Double totP39;
@Column(name = "totP40")
@Column(name = "tot_tp_40")
private Double totP40;
@Column(name = "totP41")
@Column(name = "tot_tp_41")
private Double totP41;
@Column(name = "totP42")
@Column(name = "tot_tp_42")
private Double totP42;
@Column(name = "totP43")
@Column(name = "tot_tp_43")
private Double totP43;
@Column(name = "totP44")
@Column(name = "tot_tp_44")
private Double totP44;
@Column(name = "totP45")
@Column(name = "tot_tp_45")
private Double totP45;
@Column(name = "totP46")
@Column(name = "tot_tp_46")
private Double totP46;
@Column(name = "totP47")
@Column(name = "tot_tp_47")
private Double totP47;
@Column(name = "totP49")
@Column(name = "tot_tp_49")
private Double totP48;
@Column(name = "totP49")
@Column(name = "tot_tp_49")
private Double totP49;
@Column(name = "totP50")
@Column(name = "tot_tp_50")
private Double totP50;

View File

@@ -208,107 +208,107 @@ public class DataHarmpowerQ {
private Double q50;
//在线监测添加的字段
@Column(name = "totQ")
@Column(name = "tot_q")
private Double totQ;
@Column(name = "totQ1")
@Column(name = "tot_tq_1")
private Double totQ1;
@Column(name = "totQ2")
@Column(name = "tot_tq_2")
private Double totQ2;
@Column(name = "totQ3")
@Column(name = "tot_tq_3")
private Double totQ3;
@Column(name = "totQ4")
@Column(name = "tot_tq_4")
private Double totQ4;
@Column(name = "totQ5")
@Column(name = "tot_tq_5")
private Double totQ5;
@Column(name = "totQ6")
@Column(name = "tot_tq_6")
private Double totQ6;
@Column(name = "totQ7")
@Column(name = "tot_tq_7")
private Double totQ7;
@Column(name = "totQ8")
@Column(name = "tot_tq_8")
private Double totQ8;
@Column(name = "totQ9")
@Column(name = "tot_tq_9")
private Double totQ9;
@Column(name = "totQ10")
@Column(name = "tot_tq_10")
private Double totQ10;
@Column(name = "totQ11")
@Column(name = "tot_tq_11")
private Double totQ11;
@Column(name = "totQ12")
@Column(name = "tot_tq_12")
private Double totQ12;
@Column(name = "totQ13")
@Column(name = "tot_tq_13")
private Double totQ13;
@Column(name = "totQ14")
@Column(name = "tot_tq_14")
private Double totQ14;
@Column(name = "totQ15")
@Column(name = "tot_tq_15")
private Double totQ15;
@Column(name = "totQ16")
@Column(name = "tot_tq_16")
private Double totQ16;
@Column(name = "totQ17")
@Column(name = "tot_tq_17")
private Double totQ17;
@Column(name = "totQ18")
@Column(name = "tot_tq_18")
private Double totQ18;
@Column(name = "totQ19")
@Column(name = "tot_tq_19")
private Double totQ19;
@Column(name = "totQ20")
@Column(name = "tot_tq_20")
private Double totQ20;
@Column(name = "totQ21")
@Column(name = "tot_tq_21")
private Double totQ21;
@Column(name = "totQ22")
@Column(name = "tot_tq_22")
private Double totQ22;
@Column(name = "totQ23")
@Column(name = "tot_tq_23")
private Double totQ23;
@Column(name = "totQ24")
@Column(name = "tot_tq_24")
private Double totQ24;
@Column(name = "totQ25")
@Column(name = "tot_tq_25")
private Double totQ25;
@Column(name = "totQ26")
@Column(name = "tot_tq_26")
private Double totQ26;
@Column(name = "totQ27")
@Column(name = "tot_tq_27")
private Double totQ27;
@Column(name = "totQ28")
@Column(name = "tot_tq_28")
private Double totQ28;
@Column(name = "totQ29")
@Column(name = "tot_tq_29")
private Double totQ29;
@Column(name = "totQ30")
@Column(name = "tot_tq_30")
private Double totQ30;
@Column(name = "totQ31")
@Column(name = "tot_tq_31")
private Double totQ31;
@Column(name = "totQ32")
@Column(name = "tot_tq_32")
private Double totQ32;
@Column(name = "totQ33")
@Column(name = "tot_tq_33")
private Double totQ33;
@Column(name = "totQ34")
@Column(name = "tot_tq_34")
private Double totQ34;
@Column(name = "totQ35")
@Column(name = "tot_tq_35")
private Double totQ35;
@Column(name = "totQ36")
@Column(name = "tot_tq_36")
private Double totQ36;
@Column(name = "totQ37")
@Column(name = "tot_tq_37")
private Double totQ37;
@Column(name = "totQ38")
@Column(name = "tot_tq_38")
private Double totQ38;
@Column(name = "totQ39")
@Column(name = "tot_tq_39")
private Double totQ39;
@Column(name = "totQ40")
@Column(name = "tot_tq_40")
private Double totQ40;
@Column(name = "totQ41")
@Column(name = "tot_tq_41")
private Double totQ41;
@Column(name = "totQ42")
@Column(name = "tot_tq_42")
private Double totQ42;
@Column(name = "totQ43")
@Column(name = "tot_tq_43")
private Double totQ43;
@Column(name = "totQ44")
@Column(name = "tot_tq_44")
private Double totQ44;
@Column(name = "totQ45")
@Column(name = "tot_tq_45")
private Double totQ45;
@Column(name = "totQ46")
@Column(name = "tot_tq_46")
private Double totQ46;
@Column(name = "totQ47")
@Column(name = "tot_tq_47")
private Double totQ47;
@Column(name = "totQ49")
@Column(name = "tot_tq_49")
private Double totQ48;
@Column(name = "totQ49")
@Column(name = "tot_tq_49")
private Double totQ49;
@Column(name = "totQ50")
@Column(name = "tot_tq_50")
private Double totQ50;

View File

@@ -208,107 +208,107 @@ public class DataHarmpowerS {
private Double s50;
//在线监测添加的字段
@Column(name = "totS")
@Column(name = "tot_s")
private Double totS;
@Column(name = "totS1")
@Column(name = "tot_ts_1")
private Double totS1;
@Column(name = "totS2")
@Column(name = "tot_ts_2")
private Double totS2;
@Column(name = "totS3")
@Column(name = "tot_ts_3")
private Double totS3;
@Column(name = "totS4")
@Column(name = "tot_ts_4")
private Double totS4;
@Column(name = "totS5")
@Column(name = "tot_ts_5")
private Double totS5;
@Column(name = "totS6")
@Column(name = "tot_ts_6")
private Double totS6;
@Column(name = "totS7")
@Column(name = "tot_ts_7")
private Double totS7;
@Column(name = "totS8")
@Column(name = "tot_ts_8")
private Double totS8;
@Column(name = "totS9")
@Column(name = "tot_ts_9")
private Double totS9;
@Column(name = "totS10")
@Column(name = "tot_ts_10")
private Double totS10;
@Column(name = "totS11")
@Column(name = "tot_ts_11")
private Double totS11;
@Column(name = "totS12")
@Column(name = "tot_ts_12")
private Double totS12;
@Column(name = "totS13")
@Column(name = "tot_ts_13")
private Double totS13;
@Column(name = "totS14")
@Column(name = "tot_ts_14")
private Double totS14;
@Column(name = "totS15")
@Column(name = "tot_ts_15")
private Double totS15;
@Column(name = "totS16")
@Column(name = "tot_ts_16")
private Double totS16;
@Column(name = "totS17")
@Column(name = "tot_ts_17")
private Double totS17;
@Column(name = "totS18")
@Column(name = "tot_ts_18")
private Double totS18;
@Column(name = "totS19")
@Column(name = "tot_ts_19")
private Double totS19;
@Column(name = "totS20")
@Column(name = "tot_ts_20")
private Double totS20;
@Column(name = "totS21")
@Column(name = "tot_ts_21")
private Double totS21;
@Column(name = "totS22")
@Column(name = "tot_ts_22")
private Double totS22;
@Column(name = "totS23")
@Column(name = "tot_ts_23")
private Double totS23;
@Column(name = "totS24")
@Column(name = "tot_ts_24")
private Double totS24;
@Column(name = "totS25")
@Column(name = "tot_ts_25")
private Double totS25;
@Column(name = "totS26")
@Column(name = "tot_ts_26")
private Double totS26;
@Column(name = "totS27")
@Column(name = "tot_ts_27")
private Double totS27;
@Column(name = "totS28")
@Column(name = "tot_ts_28")
private Double totS28;
@Column(name = "totS29")
@Column(name = "tot_ts_29")
private Double totS29;
@Column(name = "totS30")
@Column(name = "tot_ts_30")
private Double totS30;
@Column(name = "totS31")
@Column(name = "tot_ts_31")
private Double totS31;
@Column(name = "totS32")
@Column(name = "tot_ts_32")
private Double totS32;
@Column(name = "totS33")
@Column(name = "tot_ts_33")
private Double totS33;
@Column(name = "totS34")
@Column(name = "tot_ts_34")
private Double totS34;
@Column(name = "totS35")
@Column(name = "tot_ts_35")
private Double totS35;
@Column(name = "totS36")
@Column(name = "tot_ts_36")
private Double totS36;
@Column(name = "totS37")
@Column(name = "tot_ts_37")
private Double totS37;
@Column(name = "totS38")
@Column(name = "tot_ts_38")
private Double totS38;
@Column(name = "totS39")
@Column(name = "tot_ts_39")
private Double totS39;
@Column(name = "totS40")
@Column(name = "tot_ts_40")
private Double totS40;
@Column(name = "totS41")
@Column(name = "tot_ts_41")
private Double totS41;
@Column(name = "totS42")
@Column(name = "tot_ts_42")
private Double totS42;
@Column(name = "totS43")
@Column(name = "tot_ts_43")
private Double totS43;
@Column(name = "totS44")
@Column(name = "tot_ts_44")
private Double totS44;
@Column(name = "totS45")
@Column(name = "tot_ts_45")
private Double totS45;
@Column(name = "totS46")
@Column(name = "tot_ts_46")
private Double totS46;
@Column(name = "totS47")
@Column(name = "tot_ts_47")
private Double totS47;
@Column(name = "totS49")
@Column(name = "tot_ts_49")
private Double totS48;
@Column(name = "totS49")
@Column(name = "tot_ts_49")
private Double totS49;
@Column(name = "totS50")
@Column(name = "tot_ts_50")
private Double totS50;

View File

@@ -20,6 +20,9 @@ public class PqsCommunicateDto {
private Integer type;
//是否更新updateTime标志数据上送更新1状态翻转不更新0
private Integer flag=0;
//是否更新、pqs_communicate标志0更新1不更新
private Integer updateCommunicateFlag=0;
}

View File

@@ -12,14 +12,16 @@ import com.njcn.device.pq.api.DeviceFeignClient;
import com.njcn.device.pq.api.LineFeignClient;
import com.njcn.device.pq.pojo.dto.DevComFlagDTO;
import com.njcn.device.pq.pojo.vo.LineDeviceStateVO;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.util.CollectionUtils;
import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.stream.Collectors;
/**
@@ -30,6 +32,7 @@ import java.util.stream.Collectors;
* @version V1.0.0
*/
@Service("RelationLnDataDealServiceImpl")
@Slf4j
public class LnDataDealServiceImpl implements LnDataDealService {
@QueryBean
private IDataFlicker dataFlickerQuery;
@@ -66,33 +69,61 @@ public class LnDataDealServiceImpl implements LnDataDealService {
private LineFeignClient lineFeignClient;
@QueryBean
private IPqsCommunicate iPqsCommunicate;
@Resource(name="asyncExecutor")
private Executor executor;
@Override
public void batchInsertion(LnDataDTO lnDataDTO) {
dataVService.batchInsertion(lnDataDTO.getDataVList());
dataHarmrateVQuery.batchInsertion(lnDataDTO.getDataHarmrateVDTOList());
dataFlickerQuery.batchInsertion(lnDataDTO.getDataFlickerDTOList());
dataFlucQuery.batchInsertion(lnDataDTO.getDataFlucDTOList());
dataHarmphasicIQuery.batchInsertion(lnDataDTO.getDataHarmphasicIDTOList());
dataHarmphasicVQuery.batchInsertion(lnDataDTO.getDataHarmphasicVDTOList());
dataHarmpowerPService.batchInsertion(lnDataDTO.getDataHarmpowerPDTOList());
dataHarmpowerQService.batchInsertion(lnDataDTO.getDataHarmpowerQDTOList());
dataHarmpowerSService.batchInsertion(lnDataDTO.getDataHarmpowerSDTOList());
dataIService.batchInsertion(lnDataDTO.getDataIDTOList());
dataInharmIService.batchInsertion(lnDataDTO.getDataInharmIDTOList());
dataInharmVService.batchInsertion(lnDataDTO.getDataInharmVDTOList());
dataPltService.batchInsertion(lnDataDTO.getDataPltDTOList());
List<CompletableFuture<Void>> futures = new ArrayList<>();
// 提交每个批量插入任务
futures.add(CompletableFuture.runAsync(() -> dataVService.batchInsertion(lnDataDTO.getDataVList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataHarmrateVQuery.batchInsertion(lnDataDTO.getDataHarmrateVDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataFlickerQuery.batchInsertion(lnDataDTO.getDataFlickerDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataFlucQuery.batchInsertion(lnDataDTO.getDataFlucDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataHarmphasicIQuery.batchInsertion(lnDataDTO.getDataHarmphasicIDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataHarmphasicVQuery.batchInsertion(lnDataDTO.getDataHarmphasicVDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataHarmpowerPService.batchInsertion(lnDataDTO.getDataHarmpowerPDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataHarmpowerQService.batchInsertion(lnDataDTO.getDataHarmpowerQDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataHarmpowerSService.batchInsertion(lnDataDTO.getDataHarmpowerSDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataIService.batchInsertion(lnDataDTO.getDataIDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataInharmIService.batchInsertion(lnDataDTO.getDataInharmIDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataInharmVService.batchInsertion(lnDataDTO.getDataInharmVDTOList()), executor));
futures.add(CompletableFuture.runAsync(() -> dataPltService.batchInsertion(lnDataDTO.getDataPltDTOList()), executor));
// 等待所有任务完成,如果任一任务抛出异常,这里会重新抛出
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
// dataVService.batchInsertion(lnDataDTO.getDataVList());
// dataHarmrateVQuery.batchInsertion(lnDataDTO.getDataHarmrateVDTOList());
// dataFlickerQuery.batchInsertion(lnDataDTO.getDataFlickerDTOList());
// dataFlucQuery.batchInsertion(lnDataDTO.getDataFlucDTOList());
// dataHarmphasicIQuery.batchInsertion(lnDataDTO.getDataHarmphasicIDTOList());
// dataHarmphasicVQuery.batchInsertion(lnDataDTO.getDataHarmphasicVDTOList());
// dataHarmpowerPService.batchInsertion(lnDataDTO.getDataHarmpowerPDTOList());
// dataHarmpowerQService.batchInsertion(lnDataDTO.getDataHarmpowerQDTOList());
// dataHarmpowerSService.batchInsertion(lnDataDTO.getDataHarmpowerSDTOList());
// dataIService.batchInsertion(lnDataDTO.getDataIDTOList());
// dataInharmIService.batchInsertion(lnDataDTO.getDataInharmIDTOList());
// dataInharmVService.batchInsertion(lnDataDTO.getDataInharmVDTOList());
// dataPltService.batchInsertion(lnDataDTO.getDataPltDTOList());
//更新mysqldevice表最新数据时间
if(!CollectionUtils.isEmpty(lnDataDTO.getDataVList())){
List<String> lineIdList = lnDataDTO.getDataVList().stream().map(DataVDTO::getLineid).distinct().collect(Collectors.toList());
List<LineDeviceStateVO> data = lineFeignClient.getAllLine(lineIdList).getData();
Map<String, String> map = data.stream().collect(Collectors.toMap(LineDeviceStateVO::getId, temp -> temp.getPids().split(",")[4]));
lnDataDTO.getDataVList().forEach(temp->{
temp.setDevId(map.get(temp.getLineid()));
});
Map<String, List<DataVDTO>> collect = lnDataDTO.getDataVList().stream().collect(Collectors.groupingBy(DataVDTO::getDevId));
List<CompletableFuture<Void>> updatefutures = new ArrayList<>();
collect.forEach((temp,dataVDTOList)->{
PqsCommunicateDto pqsCommunicateDto = new PqsCommunicateDto();
DataVDTO dataVDTO =dataVDTOList.stream().max(Comparator.comparing(DataVDTO::getTimeid)).get();
@@ -101,10 +132,12 @@ public class LnDataDealServiceImpl implements LnDataDealService {
pqsCommunicateDto.setDevId(temp);
pqsCommunicateDto.setType(1);
pqsCommunicateDto.setFlag(1);
pqsCommunicateDto.setUpdateCommunicateFlag(1);
updatefutures.add(CompletableFuture.runAsync(() -> iPqsCommunicate.insertion(pqsCommunicateDto), executor));
iPqsCommunicate.insertion(pqsCommunicateDto);
});
CompletableFuture.allOf(updatefutures.toArray(new CompletableFuture[0])).join();
}

View File

@@ -3,6 +3,8 @@ package com.njcn.dataProcess.service.impl.influxdb;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.collection.CollectionUtil;
import cn.hutool.core.util.ObjectUtil;
import com.alibaba.nacos.shaded.com.google.common.reflect.TypeToken;
import com.alibaba.nacos.shaded.com.google.gson.Gson;
import com.github.jeffreyning.mybatisplus.service.MppServiceImpl;
import com.njcn.dataProcess.dao.imapper.DataFlickerMapper;
import com.njcn.dataProcess.dao.relation.mapper.RStatDataFlickerRelationMapper;
@@ -15,11 +17,14 @@ import com.njcn.dataProcess.pojo.po.RStatDataFlickerD;
import com.njcn.dataProcess.service.IDataFlicker;
import com.njcn.dataProcess.util.TimeUtils;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.redis.utils.RedisUtil;
import lombok.RequiredArgsConstructor;
import org.apache.commons.collections4.ListUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.lang.reflect.Type;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
@@ -38,15 +43,19 @@ import java.util.stream.Collectors;
public class InfluxdbDataFlickerImpl extends MppServiceImpl<RStatDataFlickerRelationMapper, RStatDataFlickerD> implements IDataFlicker {
private final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault());
private final DataFlickerMapper dataFlickerMapper;
private static final Map<String, String> PHASE_MAPPING = new HashMap<String, String>() {{
put("AB", "A");
put("BC", "B");
put("CA", "C");
put("M", "T");
}};
@Resource
private RedisUtil redisUtil;
private static final Set<String> LINE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T")));
private static final Set<String> PHASE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T")));
@Override
@@ -229,6 +238,13 @@ public class InfluxdbDataFlickerImpl extends MppServiceImpl<RStatDataFlickerRela
*/
public List<DataFlicker> getMinuteData(LineCountEvaluateParam lineParam) {
List<DataFlicker> result = new ArrayList<>();
List<DataFlicker> data = new ArrayList<>();
//获取监测点、接线方式数据
Type type = new TypeToken<Map<String, Integer>>(){}.getType();
Map<String, Integer> map = new Gson().fromJson(
String.valueOf(redisUtil.getObjectByKey("wlLineDetail")),
type
);
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataFlicker.class);
influxQueryWrapper.regular(DataFlicker::getLineId, lineParam.getLineId())
.select(DataFlicker::getLineId)
@@ -242,13 +258,30 @@ public class InfluxdbDataFlickerImpl extends MppServiceImpl<RStatDataFlickerRela
.eq(DataFlicker::getQualityFlag, "0");
quality(result, influxQueryWrapper, lineParam);
if (CollectionUtil.isNotEmpty(result)) {
result.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
if (!Objects.isNull(map)) {
//现根据监测点分组,然后根据接线方式排除多于数据,在修改相别
Map<String, List<DataFlicker>> lineMap = result.stream().collect(Collectors.groupingBy(DataFlicker::getLineId));
lineMap.forEach((k,v)->{
if (Objects.isNull(map.get(k))) {
return;
}
Integer conType = map.get(k);
Set<String> validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES;
List<DataFlicker> result2 = v.stream().filter(item -> validPhasicTypes.contains(item.getPhasicType())).collect(Collectors.toList());
data.addAll(result2);
});
} else {
data.addAll(result);
}
if (CollectionUtil.isNotEmpty(data)) {
data.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
}
}
return result;
return data;
}
}

View File

@@ -1,6 +1,8 @@
package com.njcn.dataProcess.service.impl.influxdb;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.nacos.shaded.com.google.common.reflect.TypeToken;
import com.alibaba.nacos.shaded.com.google.gson.Gson;
import com.github.jeffreyning.mybatisplus.service.MppServiceImpl;
import com.njcn.dataProcess.dao.imapper.DataFlucMapper;
import com.njcn.dataProcess.dao.relation.mapper.RStatDataFlucRelationMapper;
@@ -13,11 +15,14 @@ import com.njcn.dataProcess.pojo.po.RStatDataFlucD;
import com.njcn.dataProcess.service.IDataFluc;
import com.njcn.dataProcess.util.TimeUtils;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.redis.utils.RedisUtil;
import lombok.RequiredArgsConstructor;
import org.apache.commons.collections4.ListUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.lang.reflect.Type;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
@@ -36,15 +41,19 @@ import java.util.stream.Collectors;
public class InfluxdbDataFlucImpl extends MppServiceImpl<RStatDataFlucRelationMapper, RStatDataFlucD> implements IDataFluc {
private final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault());
private final DataFlucMapper dataFlucMapper;
private static final Map<String, String> PHASE_MAPPING = new HashMap<String, String>() {{
put("AB", "A");
put("BC", "B");
put("CA", "C");
put("M", "T");
}};
@Resource
private RedisUtil redisUtil;
private static final Set<String> LINE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T")));
private static final Set<String> PHASE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T")));
@Override
public void batchInsertion(List<DataFlucDTO> dataFlucDTOList) {
@@ -147,6 +156,13 @@ public class InfluxdbDataFlucImpl extends MppServiceImpl<RStatDataFlucRelationMa
public List<DataFluc> getMinuteData(List<String> lineList, String startTime, String endTime, Map<String,List<String>> timeMap, Boolean dataType) {
List<DataFluc> dataList;
List<DataFluc> result = new ArrayList<>();
List<DataFluc> data = new ArrayList<>();
//获取监测点、接线方式数据
Type type = new TypeToken<Map<String, Integer>>(){}.getType();
Map<String, Integer> map = new Gson().fromJson(
String.valueOf(redisUtil.getObjectByKey("wlLineDetail")),
type
);
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataFluc.class);
influxQueryWrapper.regular(DataFluc::getLineId, lineList)
.select(DataFluc::getLineId)
@@ -195,14 +211,31 @@ public class InfluxdbDataFlucImpl extends MppServiceImpl<RStatDataFlucRelationMa
}
}
if (CollectionUtil.isNotEmpty(result)) {
result.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
if (!Objects.isNull(map)) {
//现根据监测点分组,然后根据接线方式排除多于数据,在修改相别
Map<String, List<DataFluc>> lineMap = result.stream().collect(Collectors.groupingBy(DataFluc::getLineId));
lineMap.forEach((k,v)->{
if (Objects.isNull(map.get(k))) {
return;
}
Integer conType = map.get(k);
Set<String> validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES;
List<DataFluc> result2 = v.stream().filter(item -> validPhasicTypes.contains(item.getPhasicType())).collect(Collectors.toList());
data.addAll(result2);
});
} else {
data.addAll(result);
}
if (CollectionUtil.isNotEmpty(data)) {
data.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
}
}
return result;
return data;
}
}

View File

@@ -2,6 +2,8 @@ package com.njcn.dataProcess.service.impl.influxdb;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.nacos.shaded.com.google.common.reflect.TypeToken;
import com.alibaba.nacos.shaded.com.google.gson.Gson;
import com.github.jeffreyning.mybatisplus.service.MppServiceImpl;
import com.njcn.common.utils.HarmonicTimesUtil;
import com.njcn.dataProcess.constant.InfluxDBTableConstant;
@@ -9,6 +11,7 @@ import com.njcn.dataProcess.dao.imapper.DataHarmRateVMapper;
import com.njcn.dataProcess.dao.relation.mapper.RStatDataHarmRateVRelationMapper;
import com.njcn.dataProcess.dto.DataHarmrateVDTO;
import com.njcn.dataProcess.param.LineCountEvaluateParam;
import com.njcn.dataProcess.po.influx.DataFlicker;
import com.njcn.dataProcess.po.influx.DataHarmrateV;
import com.njcn.dataProcess.po.influx.DataV;
import com.njcn.dataProcess.pojo.dto.CommonMinuteDto;
@@ -19,11 +22,14 @@ import com.njcn.dataProcess.service.IDataHarmRateV;
import com.njcn.dataProcess.util.TimeUtils;
import com.njcn.influx.constant.InfluxDbSqlConstant;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.redis.utils.RedisUtil;
import lombok.RequiredArgsConstructor;
import org.apache.commons.collections4.ListUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.lang.reflect.Type;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
@@ -36,17 +42,20 @@ import java.util.stream.Collectors;
@Service("InfluxdbDataHarmRateVImpl")
@RequiredArgsConstructor
public class InfluxdbDataHarmRateVImpl extends MppServiceImpl<RStatDataHarmRateVRelationMapper, RStatDataHarmRateVD> implements IDataHarmRateV {
private final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault());
private final DataHarmRateVMapper dataHarmRateVMapper;
private static final Map<String, String> PHASE_MAPPING = new HashMap<String, String>() {{
put("AB", "A");
put("BC", "B");
put("CA", "C");
put("M", "T");
}};
@Resource
private RedisUtil redisUtil;
private static final Set<String> LINE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T")));
private static final Set<String> PHASE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T")));
@Override
public List<DataHarmDto> getRawData(LineCountEvaluateParam lineParam) {
@@ -244,6 +253,13 @@ public class InfluxdbDataHarmRateVImpl extends MppServiceImpl<RStatDataHarmRateV
public List<DataHarmrateV> getMinuteData(LineCountEvaluateParam lineParam) {
List<DataHarmrateV> dataList;
List<DataHarmrateV> result = new ArrayList<>();
List<DataHarmrateV> data = new ArrayList<>();
//获取监测点、接线方式数据
Type type = new TypeToken<Map<String, Integer>>(){}.getType();
Map<String, Integer> map = new Gson().fromJson(
String.valueOf(redisUtil.getObjectByKey("wlLineDetail")),
type
);
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataHarmrateV.class);
influxQueryWrapper.samePrefixAndSuffix(InfluxDbSqlConstant.V, InfluxDbSqlConstant.V, HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.regular(DataHarmrateV::getLineId, lineParam.getLineId())
@@ -295,13 +311,30 @@ public class InfluxdbDataHarmRateVImpl extends MppServiceImpl<RStatDataHarmRateV
}
}
if (CollectionUtil.isNotEmpty(result)) {
result.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
if (!Objects.isNull(map)) {
//现根据监测点分组,然后根据接线方式排除多于数据,在修改相别
Map<String, List<DataHarmrateV>> lineMap = result.stream().collect(Collectors.groupingBy(DataHarmrateV::getLineId));
lineMap.forEach((k,v)->{
if (Objects.isNull(map.get(k))) {
return;
}
Integer conType = map.get(k);
Set<String> validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES;
List<DataHarmrateV> result2 = v.stream().filter(item -> validPhasicTypes.contains(item.getPhasicType())).collect(Collectors.toList());
data.addAll(result2);
});
} else {
data.addAll(result);
}
if (CollectionUtil.isNotEmpty(data)) {
data.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
}
}
return result;
return data;
}
}

View File

@@ -2,6 +2,8 @@ package com.njcn.dataProcess.service.impl.influxdb;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.nacos.shaded.com.google.common.reflect.TypeToken;
import com.alibaba.nacos.shaded.com.google.gson.Gson;
import com.github.jeffreyning.mybatisplus.service.MppServiceImpl;
import com.njcn.common.utils.HarmonicTimesUtil;
import com.njcn.dataProcess.dao.imapper.DataHarmphasicVMapper;
@@ -17,11 +19,14 @@ import com.njcn.dataProcess.service.IDataHarmphasicV;
import com.njcn.dataProcess.util.TimeUtils;
import com.njcn.influx.constant.InfluxDbSqlConstant;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.redis.utils.RedisUtil;
import lombok.RequiredArgsConstructor;
import org.apache.commons.collections4.ListUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.lang.reflect.Type;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
@@ -41,13 +46,18 @@ public class InfluxdbDataHarmphasicVImpl extends MppServiceImpl<RStatDataHarmPha
private final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault());
private final DataHarmphasicVMapper dataHarmphasicVMapper;
private static final Map<String, String> PHASE_MAPPING = new HashMap<String, String>() {{
put("AB", "A");
put("BC", "B");
put("CA", "C");
put("M", "T");
}};
@Resource
private RedisUtil redisUtil;
private static final Set<String> LINE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T")));
private static final Set<String> PHASE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T")));
@Override
@@ -205,6 +215,13 @@ public class InfluxdbDataHarmphasicVImpl extends MppServiceImpl<RStatDataHarmPha
public List<DataHarmphasicV> getMinuteData(List<String> lineList, String startTime, String endTime, Map<String,List<String>> timeMap, Boolean dataType) {
List<DataHarmphasicV> dataList;
List<DataHarmphasicV> result = new ArrayList<>();
List<DataHarmphasicV> data = new ArrayList<>();
//获取监测点、接线方式数据
Type type = new TypeToken<Map<String, Integer>>(){}.getType();
Map<String, Integer> map = new Gson().fromJson(
String.valueOf(redisUtil.getObjectByKey("wlLineDetail")),
type
);
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataHarmphasicV.class);
influxQueryWrapper.samePrefixAndSuffix(InfluxDbSqlConstant.V, InfluxDbSqlConstant.V, HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.regular(DataHarmphasicV::getLineId, lineList)
@@ -253,13 +270,30 @@ public class InfluxdbDataHarmphasicVImpl extends MppServiceImpl<RStatDataHarmPha
}
}
if (CollectionUtil.isNotEmpty(result)) {
result.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
if (!Objects.isNull(map)) {
//现根据监测点分组,然后根据接线方式排除多于数据,在修改相别
Map<String, List<DataHarmphasicV>> lineMap = result.stream().collect(Collectors.groupingBy(DataHarmphasicV::getLineId));
lineMap.forEach((k,v)->{
if (Objects.isNull(map.get(k))) {
return;
}
Integer conType = map.get(k);
Set<String> validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES;
List<DataHarmphasicV> result2 = v.stream().filter(item -> validPhasicTypes.contains(item.getPhasicType())).collect(Collectors.toList());
data.addAll(result2);
});
} else {
data.addAll(result);
}
if (CollectionUtil.isNotEmpty(data)) {
data.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
}
}
return result;
return data;
}
}

View File

@@ -192,6 +192,7 @@ public class InfluxdbDataHarmpowerPImpl extends MppServiceImpl<RStatDataHarmPowe
List<DataHarmpowerP> result = new ArrayList<>();
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataHarmpowerP.class);
influxQueryWrapper.samePrefixAndSuffix(InfluxDbSqlConstant.P, InfluxDbSqlConstant.P, HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.samePrefixAndSuffix("tot_tp_", "tot_tp_", HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.regular(DataHarmpowerP::getLineId, lineList)
.select(DataHarmpowerP::getLineId)
.select(DataHarmpowerP::getPhasicType)
@@ -199,6 +200,9 @@ public class InfluxdbDataHarmpowerPImpl extends MppServiceImpl<RStatDataHarmPowe
.select(DataHarmpowerP::getP)
.select(DataHarmpowerP::getDf)
.select(DataHarmpowerP::getPf)
.select(DataHarmpowerP::getTotP)
.select(DataHarmpowerP::getTotDf)
.select(DataHarmpowerP::getTotPf)
.select(DataHarmpowerP::getQualityFlag)
.select(DataHarmpowerP::getAbnormalFlag)
.between(DataHarmpowerP::getTime, startTime, endTime)

View File

@@ -105,7 +105,6 @@ public class InfluxdbDataHarmpowerQImpl extends MppServiceImpl<RStatDataHarmPowe
valueTypeMap.forEach((valueType,valueTypeList)->{
CommonMinuteDto.ValueType value = new CommonMinuteDto.ValueType();
value.setValueType(valueType);
List<List<Double>> lists;
if (Objects.equals(phasicType, "T") && Objects.equals(lineParam.getType(), 2)) {
lists = extractDataLists(valueTypeList, "Tot");
@@ -188,11 +187,13 @@ public class InfluxdbDataHarmpowerQImpl extends MppServiceImpl<RStatDataHarmPowe
List<DataHarmpowerQ> result = new ArrayList<>();
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataHarmpowerQ.class);
influxQueryWrapper.samePrefixAndSuffix(InfluxDbSqlConstant.Q, InfluxDbSqlConstant.Q, HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.samePrefixAndSuffix("tot_tq_", "tot_tq_", HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.regular(DataHarmpowerQ::getLineId, lineList)
.select(DataHarmpowerQ::getLineId)
.select(DataHarmpowerQ::getPhasicType)
.select(DataHarmpowerQ::getValueType)
.select(DataHarmpowerQ::getQ)
.select(DataHarmpowerQ::getTotQ)
.select(DataHarmpowerQ::getQualityFlag)
.select(DataHarmpowerQ::getAbnormalFlag)
.between(DataHarmpowerQ::getTime, startTime, endTime)

View File

@@ -187,11 +187,13 @@ public class InfluxdbDataHarmpowerSImpl extends MppServiceImpl<RStatDataHarmPowe
List<DataHarmpowerS> result = new ArrayList<>();
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataHarmpowerS.class);
influxQueryWrapper.samePrefixAndSuffix(InfluxDbSqlConstant.S, InfluxDbSqlConstant.S, HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.samePrefixAndSuffix("tot_ts_", "tot_ts_", HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.regular(DataHarmpowerS::getLineId, lineList)
.select(DataHarmpowerS::getLineId)
.select(DataHarmpowerS::getPhasicType)
.select(DataHarmpowerS::getValueType)
.select(DataHarmpowerS::getS)
.select(DataHarmpowerS::getTotS)
.select(DataHarmpowerS::getQualityFlag)
.select(DataHarmpowerS::getAbnormalFlag)
.between(DataHarmpowerS::getTime, startTime, endTime)

View File

@@ -2,6 +2,8 @@ package com.njcn.dataProcess.service.impl.influxdb;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.nacos.shaded.com.google.common.reflect.TypeToken;
import com.alibaba.nacos.shaded.com.google.gson.Gson;
import com.github.jeffreyning.mybatisplus.service.MppServiceImpl;
import com.njcn.dataProcess.dao.imapper.DataPltMapper;
import com.njcn.dataProcess.dao.relation.mapper.RStatDataPltRelationMapper;
@@ -15,11 +17,14 @@ import com.njcn.dataProcess.pojo.po.RStatDataPltD;
import com.njcn.dataProcess.service.IDataPlt;
import com.njcn.dataProcess.util.TimeUtils;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.redis.utils.RedisUtil;
import lombok.RequiredArgsConstructor;
import org.apache.commons.collections4.ListUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.lang.reflect.Type;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
@@ -39,13 +44,18 @@ public class InfluxdbDataPltImpl extends MppServiceImpl<RStatDataPltRelationMapp
private final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault());
private final DataPltMapper dataPltMapper;
private static final Map<String, String> PHASE_MAPPING = new HashMap<String, String>() {{
put("AB", "A");
put("BC", "B");
put("CA", "C");
put("M", "T");
}};
@Resource
private RedisUtil redisUtil;
private static final Set<String> LINE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T")));
private static final Set<String> PHASE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T")));
@Override
public void batchInsertion(List<DataPltDTO> dataPltDTOList) {
@@ -152,6 +162,13 @@ public class InfluxdbDataPltImpl extends MppServiceImpl<RStatDataPltRelationMapp
public List<DataPlt> getMinuteDataPlt(LineCountEvaluateParam lineParam) {
List<DataPlt> dataList;
List<DataPlt> result = new ArrayList<>();
List<DataPlt> data = new ArrayList<>();
//获取监测点、接线方式数据
Type type = new TypeToken<Map<String, Integer>>(){}.getType();
Map<String, Integer> map = new Gson().fromJson(
String.valueOf(redisUtil.getObjectByKey("wlLineDetail")),
type
);
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataPlt.class);
influxQueryWrapper.regular(DataPlt::getLineId, lineParam.getLineId())
.select(DataPlt::getLineId)
@@ -202,13 +219,30 @@ public class InfluxdbDataPltImpl extends MppServiceImpl<RStatDataPltRelationMapp
}
}
if (CollectionUtil.isNotEmpty(result)) {
result.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
if (!Objects.isNull(map)) {
//现根据监测点分组,然后根据接线方式排除多于数据,在修改相别
Map<String, List<DataPlt>> lineMap = result.stream().collect(Collectors.groupingBy(DataPlt::getLineId));
lineMap.forEach((k,v)->{
if (Objects.isNull(map.get(k))) {
return;
}
Integer conType = map.get(k);
Set<String> validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES;
List<DataPlt> result2 = v.stream().filter(item -> validPhasicTypes.contains(item.getPhasicType())).collect(Collectors.toList());
data.addAll(result2);
});
} else {
data.addAll(result);
}
if (CollectionUtil.isNotEmpty(data)) {
data.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
}
}
return result;
return data;
}
}

View File

@@ -4,6 +4,8 @@ import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.collection.CollectionUtil;
import cn.hutool.core.collection.ListUtil;
import cn.hutool.core.util.ObjectUtil;
import com.alibaba.nacos.shaded.com.google.common.reflect.TypeToken;
import com.alibaba.nacos.shaded.com.google.gson.Gson;
import com.github.jeffreyning.mybatisplus.service.MppServiceImpl;
import com.njcn.common.utils.HarmonicTimesUtil;
import com.njcn.dataProcess.constant.InfluxDBTableConstant;
@@ -24,16 +26,19 @@ import com.njcn.dataProcess.service.IDataV;
import com.njcn.dataProcess.util.TimeUtils;
import com.njcn.influx.constant.InfluxDbSqlConstant;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.redis.utils.RedisUtil;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections4.ListUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.lang.reflect.Type;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
@@ -46,16 +51,20 @@ import java.util.stream.Collectors;
public class InfluxdbDataVImpl extends MppServiceImpl<RStatDataVRelationMapper, RStatDataVD> implements IDataV {
private final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault());
private static final Map<String, String> PHASE_MAPPING = new HashMap<String, String>() {{
put("AB", "A");
put("BC", "B");
put("CA", "C");
put("M", "T");
}};
@Resource
private DataVMapper dataVMapper;
@Resource
private RedisUtil redisUtil;
private static final Set<String> LINE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T")));
private static final Set<String> PHASE_VOLTAGE_TYPES =
Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T")));
/**
* 注意influxdb不推荐采用in函数的方式批量查询监测点的数据效率很低容易造成崩溃故每次单测点查询
@@ -437,6 +446,13 @@ public class InfluxdbDataVImpl extends MppServiceImpl<RStatDataVRelationMapper,
*/
public List<DataV> getMinuteDataV(LineCountEvaluateParam lineParam) {
List<DataV> result = new ArrayList<>();
List<DataV> data = new ArrayList<>();
//获取监测点、接线方式数据
Type type = new TypeToken<Map<String, Integer>>(){}.getType();
Map<String, Integer> map = new Gson().fromJson(
String.valueOf(redisUtil.getObjectByKey("wlLineDetail")),
type
);
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataV.class);
influxQueryWrapper.samePrefixAndSuffix(InfluxDbSqlConstant.V, InfluxDbSqlConstant.V, HarmonicTimesUtil.harmonicTimesList(1, 50, 1));
influxQueryWrapper.regular(DataV::getLineId, lineParam.getLineId())
@@ -463,14 +479,74 @@ public class InfluxdbDataVImpl extends MppServiceImpl<RStatDataVRelationMapper,
}
quality(result, influxQueryWrapper, lineParam);
if (CollectionUtil.isNotEmpty(result)) {
result.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
if (!Objects.isNull(map)) {
//现根据监测点分组,然后根据接线方式排除多于数据,在修改相别
Map<String, List<DataV>> lineMap = result.stream().collect(Collectors.groupingBy(DataV::getLineId));
lineMap.forEach((k,v)->{
if (Objects.isNull(map.get(k))) {
return;
}
//这边需要特殊处理下,将线电压数据赋值
Map<String, DataV> lineVoltageIndex = v.stream()
.filter(d -> PHASE_MAPPING.containsKey(d.getPhasicType()))
.filter(d -> d.getRmsLvr() != null)
.collect(Collectors.toMap(
d -> buildKey(d.getTime(), d.getValueType(), d.getPhasicType()),
Function.identity(),
(existing, replacement) -> existing
));
v.stream()
.filter(d -> PHASE_VOLTAGE_TYPES.contains(d.getPhasicType()))
.forEach(phaseData -> {
// 根据当前相电压反查对应的线电压相别
String targetLinePhasic = getReverseLinePhasic(phaseData.getPhasicType());
if (targetLinePhasic == null) {
return;
}
String key = buildKey(phaseData.getTime(), phaseData.getValueType(), targetLinePhasic);
DataV matchedLineData = lineVoltageIndex.get(key);
if (matchedLineData != null && matchedLineData.getRmsLvr() != null) {
phaseData.setRmsLvr(matchedLineData.getRmsLvr());
}
});
Integer conType = map.get(k);
Set<String> validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES;
List<DataV> result2 = v.stream().filter(item -> validPhasicTypes.contains(item.getPhasicType())).collect(Collectors.toList());
data.addAll(result2);
});
} else {
data.addAll(result);
}
if (CollectionUtil.isNotEmpty(data)) {
data.forEach(item -> {
String newType = PHASE_MAPPING.get(item.getPhasicType());
if (newType != null) {
item.setPhasicType(newType);
}
});
}
}
return data;
}
private static String buildKey(Object time, Object valueType, Object phasicType) {
return time + "|" + valueType + "|" + phasicType;
}
private static String getReverseLinePhasic(String phaseType) {
if (phaseType == null) {
return null;
}
switch (phaseType) {
case "A":
return "AB";
case "B":
return "BC";
case "C":
return "CA";
default:
return null;
}
return result;
}
private void quality(List<DataV> result, InfluxQueryWrapper influxQueryWrapper, LineCountEvaluateParam lineParam) {

View File

@@ -110,23 +110,32 @@ public class InfluxdbPqsCommunicateImpl implements IPqsCommunicate {
@Override
public void insertion(PqsCommunicateDto pqsCommunicateDto) {
// log.info("进出Influxdb实现类");
//获取最新一条数据
PqsCommunicate dto = new PqsCommunicate();
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(PqsCommunicate.class);
influxQueryWrapper.eq(PqsCommunicate::getDevId,pqsCommunicateDto.getDevId()).timeDesc().limit(1);
List<PqsCommunicate> pqsCommunicates = pqsCommunicateMapper.selectByQueryWrapper(influxQueryWrapper);
if(Objects.equals(pqsCommunicateDto.getUpdateCommunicateFlag(),0)){
PqsCommunicate dto = new PqsCommunicate();
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(PqsCommunicate.class);
influxQueryWrapper.eq(PqsCommunicate::getDevId,pqsCommunicateDto.getDevId()).timeDesc().limit(1);
List<PqsCommunicate> pqsCommunicates = pqsCommunicateMapper.selectByQueryWrapper(influxQueryWrapper);
PqsCommunicate pqsCommunicate = new PqsCommunicate();
pqsCommunicate.setTime(LocalDateTime.parse(pqsCommunicateDto.getTime(), DATE_TIME_FORMATTER).atZone(ZoneId.systemDefault()).toInstant());
pqsCommunicate.setDevId(pqsCommunicateDto.getDevId());
pqsCommunicate.setType(pqsCommunicateDto.getType());
//如果不存数据或者状态不一样则插入数据
//可能存在掉线后最后一组数据还未入库,添加时间判断
if(CollectionUtils.isEmpty(pqsCommunicates)|| (!Objects.equals( pqsCommunicates.get(0).getType(),pqsCommunicateDto.getType())&&pqsCommunicates.get(0).getTime().isBefore(pqsCommunicate.getTime()))){
pqsCommunicateMapper.insertOne(pqsCommunicate);
}
}
PqsCommunicate pqsCommunicate = new PqsCommunicate();
pqsCommunicate.setTime(LocalDateTime.parse(pqsCommunicateDto.getTime(), DATE_TIME_FORMATTER).atZone(ZoneId.systemDefault()).toInstant());
pqsCommunicate.setDevId(pqsCommunicateDto.getDevId());
pqsCommunicate.setType(pqsCommunicateDto.getType());
//如果不存数据或者状态不一样则插入数据
//可能存在掉线后最后一组数据还未入库,添加时间判断
if(CollectionUtils.isEmpty(pqsCommunicates)|| (!Objects.equals( pqsCommunicates.get(0).getType(),pqsCommunicateDto.getType())&&pqsCommunicates.get(0).getTime().isBefore(pqsCommunicate.getTime()))){
pqsCommunicateMapper.insertOne(pqsCommunicate);
}
//更新mysql数据
DevComFlagDTO devComFlagDTO = new DevComFlagDTO();
devComFlagDTO.setId(pqsCommunicateDto.getDevId());
@@ -136,6 +145,8 @@ public class InfluxdbPqsCommunicateImpl implements IPqsCommunicate {
devComFlagDTO.setStatus(pqsCommunicateDto.getType());
deviceFeignClient.updateDevComFlag(devComFlagDTO);
}

View File

@@ -31,6 +31,8 @@ public class DeviceRebootMessage {
private String ip;
//终端型号
private String devType;
//前置类型
private String comType;
//挂载单位
private String org_name;
//组织名称

View File

@@ -100,12 +100,34 @@
<artifactId>spring-boot-starter-websocket</artifactId>
<version>2.7.12</version>
</dependency>
<dependency>
<groupId>com.njcn</groupId>
<artifactId>mq-spring-boot-starter</artifactId>
<version>1.0.0-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<version>2.7.12</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.13.2</version>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<finalName>messageboot</finalName>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>2.22.2</version>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>

View File

@@ -39,6 +39,8 @@ import java.util.concurrent.TimeUnit;
* @createTime 2023/8/11 15:32
*/
@Component
@org.springframework.boot.autoconfigure.condition.ConditionalOnProperty(
name = "mq.type", havingValue = "rocketmq", matchIfMissing = true)
@RocketMQMessageListener(
topic = "LN_Topic",
consumerGroup = "ln_consumer",

View File

@@ -0,0 +1,23 @@
package com.njcn.message.mq;
import com.njcn.message.messagedto.MessageDataDTO;
import com.njcn.mq.annotation.MqListener;
import com.njcn.stat.api.MessAnalysisFeignClient;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.Collections;
@Component
@ConditionalOnProperty(name = "mq.type", havingValue = "redis-stream")
public class FrontDataMqListener {
@Resource
private MessAnalysisFeignClient messAnalysisFeignClient;
@MqListener(topic = "LN_Topic", group = "ln_consumer")
public void onFrontData(MessageDataDTO msg) {
messAnalysisFeignClient.analysis(Collections.singletonList(msg));
}
}

View File

@@ -0,0 +1,19 @@
package com.njcn.message.mq;
import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.RecallMessage;
import com.njcn.mq.core.MqTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class RecallMqSender {
@Resource
private MqTemplate mqTemplate;
public void send(RecallMessage message, String nodeId) {
mqTemplate.send(nodeId + "_" + BusinessTopic.RECALL_TOPIC, message);
}
}

View File

@@ -0,0 +1,210 @@
package com.njcn.message.mq;
import com.alibaba.fastjson.JSON;
import com.njcn.message.message.RecallMessage;
import com.njcn.message.messagedto.MessageDataDTO;
import com.njcn.middle.stream.autoconfig.RedisStreamAutoConfiguration;
import com.njcn.mq.autoconfig.MqCoreAutoConfiguration;
import com.njcn.mq.container.MqListenerRegistry;
import com.njcn.mq.core.MqTemplate;
import com.njcn.mq.driver.redis.RedisStreamMqDriver;
import com.njcn.mq.driver.redis.RedisStreamMqDriverAutoConfiguration;
import com.njcn.mq.driver.rocketmq.RocketMqDriver;
import com.njcn.mq.driver.rocketmq.RocketMqDriverAutoConfiguration;
import com.njcn.mq.spi.MqDriver;
import com.njcn.stat.api.MessAnalysisFeignClient;
import org.apache.rocketmq.spring.autoconfigure.RocketMQProperties;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.Test;
import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.data.redis.connection.stream.MapRecord;
import org.springframework.data.redis.connection.stream.ReadOffset;
import org.springframework.data.redis.connection.stream.StreamOffset;
import org.springframework.data.redis.connection.stream.StreamReadOptions;
import org.springframework.data.redis.connection.stream.StreamRecords;
import org.springframework.data.redis.core.StringRedisTemplate;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
/**
* MQ starter 切片接入验证测试ApplicationContextRunner不启动全应用
* 真实 Redis 要求localhost:6379docker redis-stream-it仅上/下行收发用例需要,
* 纯装配用例不需要Lettuce 懒连接,不触发实际连接)。
*
* <p>装配一律用 {@code withConfiguration(AutoConfigurations.of(...))},让 Spring 按
* <b>真实 autoconfig 顺序</b>AutoConfigurationSorter处理与生产环境一致。以此
* 验证 {@link MqCoreAutoConfiguration} 的 {@code @ConditionalOnBean(MqDriver)} 能在
* 真实顺序下拿到 driver bean —— 这依赖 driver autoconfig 上的
* {@code @AutoConfigureBefore(MqCoreAutoConfiguration.class)}。若缺该保障,按类名字母序
* {@code com.njcn.mq.autoconfig.*} 先于 {@code com.njcn.mq.driver.*}core 会在 driver
* bean 注册前求值条件失败,导致 MqTemplate / MqListenerRegistry 不装配。
*/
public class MqSliceWiringTest {
private static boolean redisAvailable() {
try {
Socket s = new Socket();
s.connect(new InetSocketAddress("localhost", 6379), 500);
s.close();
return true;
} catch (Exception e) {
return false;
}
}
/** redis-stream 模式:真实 autoconfig 顺序装 5 个 autoconfig + 业务 bean上/下行用)。 */
private ApplicationContextRunner redisStreamRunner() {
return new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(
RedisAutoConfiguration.class,
RedisStreamAutoConfiguration.class,
RedisStreamMqDriverAutoConfiguration.class,
RocketMqDriverAutoConfiguration.class,
MqCoreAutoConfiguration.class))
.withUserConfiguration(FrontDataMqListener.class, RecallMqSender.class)
.withBean(MessAnalysisFeignClient.class,
() -> mock(MessAnalysisFeignClient.class))
.withPropertyValues(
"mq.type=redis-stream",
"spring.redis.host=localhost",
"spring.redis.port=6379");
}
// ------------------------------------------------------------------
// 1. 装配mq.type 二选一 + core bean 在真实 autoconfig 顺序下装配
// ------------------------------------------------------------------
@Test
void testRedisStreamDriverSelected() {
// 纯装配检查,不依赖真实 RedisLettuce 懒连接)。
new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(
RedisAutoConfiguration.class,
RedisStreamAutoConfiguration.class,
RedisStreamMqDriverAutoConfiguration.class,
RocketMqDriverAutoConfiguration.class,
MqCoreAutoConfiguration.class))
.withPropertyValues(
"mq.type=redis-stream",
"spring.redis.host=localhost",
"spring.redis.port=6379")
.run(ctx -> {
assertThat(ctx).hasSingleBean(MqDriver.class);
assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RedisStreamMqDriver.class);
// 真实顺序下 core 的 @ConditionalOnBean(MqDriver) 必须命中
assertThat(ctx).hasSingleBean(MqTemplate.class);
assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
});
}
@Test
void testRocketMqDriverSelected() {
// rocket 测试不依赖真实 Redis不装 RedisAutoConfigurationredis driver 因 mq.type 不匹配而 inactive。
new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(
RedisStreamMqDriverAutoConfiguration.class,
RocketMqDriverAutoConfiguration.class,
MqCoreAutoConfiguration.class))
.withBean(RocketMQTemplate.class, () -> mock(RocketMQTemplate.class))
.withBean(RocketMQProperties.class, RocketMQProperties::new)
.withPropertyValues("mq.type=rocketmq")
.run(ctx -> {
assertThat(ctx).hasSingleBean(MqDriver.class);
assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RocketMqDriver.class);
assertThat(ctx).hasSingleBean(MqTemplate.class);
assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
});
}
// ------------------------------------------------------------------
// 2. 上行(需要真实 Redis
// 向 LN_Topic stream XADD 一条消息,验证 FrontDataMqListener 触发了
// messAnalysisFeignClient.analysis(...)
// ------------------------------------------------------------------
@Test
void testUpstreamFrontData() {
Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
redisStreamRunner().run(ctx -> {
StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
MessAnalysisFeignClient mockClient = ctx.getBean(MessAnalysisFeignClient.class);
// 构造 RedisStreamMqDriver 消费时期望的字段格式
// enc=plain → body 即为原始 JSON参见 MessageCodec.encodeBody(json, false)
// stream key = StreamKeyBuilder.buildStream("LN_Topic") = "LN_Topic"
// (默认 envIsolation=false无 env 前缀)
MessageDataDTO dto = new MessageDataDTO();
dto.setDataType(1);
String json = JSON.toJSONString(dto);
Map<String, String> fields = new HashMap<String, String>();
fields.put("key", "test-key-upstream-1");
fields.put("tag", "");
fields.put("enc", "plain");
fields.put("body", json);
fields.put("ts", String.valueOf(System.currentTimeMillis()));
// XADD 到 LN_TopicMqListenerRegistry.start() 已在 context 启动时订阅该 stream
redis.opsForStream().add(
StreamRecords.newRecord().in("LN_Topic").ofMap(fields)
);
// 等待消费线程处理blockMs 默认 2000mstimeout=5000ms 足够)
verify(mockClient, timeout(5000)).analysis(anyList());
});
}
// ------------------------------------------------------------------
// 3. 下行(需要真实 Redis
// 通过 RecallMqSender 发送 RecallMessage验证 stream 内容正确
// ------------------------------------------------------------------
@Test
void testDownstreamRecall() {
Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
redisStreamRunner().run(ctx -> {
StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
RecallMqSender sender = ctx.getBean(RecallMqSender.class);
// 清理,确保只有本次发送的消息
redis.delete("n1_recall_Topic");
RecallMessage msg = new RecallMessage();
msg.setGuid("g1");
// MqTemplate.send("n1_recall_Topic", msg) → RedisStreamMqDriver XADD
// stream key = StreamKeyBuilder.buildStream("n1_recall_Topic") = "n1_recall_Topic"
sender.send(msg, "n1");
Long size = redis.opsForStream().size("n1_recall_Topic");
assertThat(size).isGreaterThanOrEqualTo(1L);
@SuppressWarnings("unchecked")
List<MapRecord<String, Object, Object>> records =
redis.opsForStream().read(
StreamReadOptions.empty().count(1L),
StreamOffset.create("n1_recall_Topic", ReadOffset.from("0-0"))
);
assertThat(records).isNotEmpty();
Map<Object, Object> firstFields = records.get(0).getValue();
// DefaultMqTemplate.send(topic, payload) 使用 encodeBody(json, false) → enc="plain"
// 因此 body 即为原始 JSON无需 gzip 解压
String body = (String) firstFields.get("body");
RecallMessage result = JSON.parseObject(body, RecallMessage.class);
assertThat(result.getGuid()).isEqualTo("g1");
});
}
}