diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/ExecutionCenter.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/ExecutionCenter.java index 2c60b7c..e01a5fe 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/ExecutionCenter.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/ExecutionCenter.java @@ -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 csLineDTOList = csLineFeignClient.getAllLineDetail().getData(); + List lineIdList = new ArrayList<>(); + Map 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()) { diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/service/line/IDataCleanService.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/service/line/IDataCleanService.java index aa4742a..2c9f6ca 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/service/line/IDataCleanService.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/service/line/IDataCleanService.java @@ -2,7 +2,6 @@ package com.njcn.algorithm.service.line; import com.njcn.algorithm.pojo.bo.CalculatedParam; -import java.util.concurrent.CompletableFuture; /** * @author xy diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataCleanServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataCleanServiceImpl.java index eb1ef64..a2c6cac 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataCleanServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataCleanServiceImpl.java @@ -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 listOfString = (List) (List) calculatedParam.getIdList(); List 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> 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)); } diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataComAssServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataComAssServiceImpl.java index a5f9cdd..ab3155b 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataComAssServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DataComAssServiceImpl.java @@ -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 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)) { diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DayDataServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DayDataServiceImpl.java index 65ad1d4..d7e7bd9 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DayDataServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/DayDataServiceImpl.java @@ -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 result = new ArrayList<>(); //远程接口获取分钟数据 LineCountEvaluateParam lineParam = new LineCountEvaluateParam(); @@ -91,7 +93,6 @@ public class DayDataServiceImpl implements IDayDataService { //获取原始数据 List partList = dataVFeignClient.getBaseData(lineParam).getData(); if (CollUtil.isNotEmpty(partList)) { - logger.info("{}dataV集合大小为>>>>>>>>>>>>{}",lineParam.getStartTime(),MemorySizeUtil.getObjectSize(partList)); partList.forEach(item->{ //相别 List 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 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 valueList = Arrays.asList("AVG","MAX","MIN","CP95"); List 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 valueList = Arrays.asList("AVG","MAX","MIN","CP95"); List 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 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 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 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 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 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 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 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 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 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 result = new ArrayList<>(); List valueList = Arrays.asList("AVG","MAX","MIN","CP95"); //远程接口获取分钟数据 diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/FlowAsyncServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/FlowAsyncServiceImpl.java index 51acadd..ce276f9 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/FlowAsyncServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/FlowAsyncServiceImpl.java @@ -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(); } diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataCrossingServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataCrossingServiceImpl.java index 1ca4af5..38ee5ba 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataCrossingServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataCrossingServiceImpl.java @@ -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 overLimitMap = overLimitList.stream().collect(Collectors.toMap(Overlimit::getId, Function.identity())); //以100个监测点分片处理 - List> pendingIds = ListUtils.partition(lineIds, 1); - ArrayList phase = ListUtil.toList(PhaseType.PHASE_A, PhaseType.PHASE_B, PhaseType.PHASE_C); MemorySizeUtil.getNowMemory(); + List> pendingIds = ListUtils.partition(lineIds, 20); + ArrayList phase = ListUtil.toList(PhaseType.PHASE_A, PhaseType.PHASE_B, PhaseType.PHASE_C); - List> futures = new ArrayList<>(); - for (int i = 0; i < pendingIds.size(); i++) { - logger.info(calculatedParam.getDataDate()+" 总分区数据:" + pendingIds.size() + "=====》当前第"+(i + 1)+"小分区"); - List list = pendingIds.get(i); - // 获取Future - CompletableFuture future = dataLimitRateAsync.lineDataRate( - calculatedParam.getDataDate(), - list, - phase, - overLimitMap, - pendingIds.size(), - (i + 1), - lineParam.getType() - ); - - futures.add(future); + for (List pendingId : pendingIds) { + List> futures = new ArrayList<>(); + for (String s : pendingId) { + CompletableFuture 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())); diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataIntegrityServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataIntegrityServiceImpl.java index 8c67c98..6136050 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataIntegrityServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataIntegrityServiceImpl.java @@ -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 calculatedParam) { - logger.info("{},integrity表转r_stat_integrity_d算法开始=====》", LocalDateTime.now()); + MemorySizeUtil.getNowMemory(); + logger.info("{},{}integrity表转r_stat_integrity_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate()); List poList = new ArrayList<>(); List lineIds = calculatedParam.getIdList(); String beginDay = LocalDateTimeUtil.format( diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataLimitRateAsyncImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataLimitRateAsyncImpl.java index 73d0126..5abb704 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataLimitRateAsyncImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataLimitRateAsyncImpl.java @@ -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); diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataOnlineRateServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataOnlineRateServiceImpl.java index c5a6a2b..e13b950 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataOnlineRateServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IDataOnlineRateServiceImpl.java @@ -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())); diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IEventDetailServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IEventDetailServiceImpl.java index 0c2b72e..f2dafe3 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IEventDetailServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/IEventDetailServiceImpl.java @@ -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 calculatedParam) { - logger.info("{},r_mp_event_detail_d算法开始=====》", LocalDateTime.now()); + MemorySizeUtil.getNowMemory(); + logger.info("{},{}r_mp_event_detail_d算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate()); List result = new ArrayList<>(); //远程接口获取分钟数据 LineCountEvaluateParam lineParam = new LineCountEvaluateParam(); @@ -106,7 +108,8 @@ public class IEventDetailServiceImpl implements IEventDetailService { @Override public void dataMonthHandle(CalculatedParam calculatedParam) { - logger.info("{},r_mp_event_detail_m算法开始=====》", LocalDateTime.now()); + MemorySizeUtil.getNowMemory(); + logger.info("{},{}r_mp_event_detail_m算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate()); List result = new ArrayList<>(); //远程接口获取分钟数据 LineCountEvaluateParam lineParam = new LineCountEvaluateParam(); @@ -139,7 +142,8 @@ public class IEventDetailServiceImpl implements IEventDetailService { @Override public void dataQuarterHandle(CalculatedParam calculatedParam) { - logger.info("{},r_mp_event_detail_q算法开始=====》", LocalDateTime.now()); + MemorySizeUtil.getNowMemory(); + logger.info("{},{}r_mp_event_detail_q算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate()); List result = new ArrayList<>(); //远程接口获取分钟数据 LineCountEvaluateParam lineParam = new LineCountEvaluateParam(); @@ -172,7 +176,8 @@ public class IEventDetailServiceImpl implements IEventDetailService { @Override public void dataYearHandle(CalculatedParam calculatedParam) { - logger.info("{},r_mp_event_detail_y算法开始=====》", LocalDateTime.now()); + MemorySizeUtil.getNowMemory(); + logger.info("{},{}r_mp_event_detail_y算法开始=====》", LocalDateTime.now(),calculatedParam.getDataDate()); List result = new ArrayList<>(); //远程接口获取分钟数据 LineCountEvaluateParam lineParam = new LineCountEvaluateParam(); diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/PollutionServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/PollutionServiceImpl.java index 4092d68..ac0be54 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/PollutionServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/PollutionServiceImpl.java @@ -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; diff --git a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/SpecialAnalysisServiceImpl.java b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/SpecialAnalysisServiceImpl.java index 85ceb40..9255510 100644 --- a/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/SpecialAnalysisServiceImpl.java +++ b/algorithm/algorithm-boot/src/main/java/com/njcn/algorithm/serviceimpl/line/SpecialAnalysisServiceImpl.java @@ -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; diff --git a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/api/SpThroughFeignClient.java b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/api/SpThroughFeignClient.java index d438d19..bcd44e3 100644 --- a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/api/SpThroughFeignClient.java +++ b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/api/SpThroughFeignClient.java @@ -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; diff --git a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerP.java b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerP.java index 49d9cc3..1e75132 100644 --- a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerP.java +++ b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerP.java @@ -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; diff --git a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerQ.java b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerQ.java index beaeb33..5ca76a7 100644 --- a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerQ.java +++ b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerQ.java @@ -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; diff --git a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerS.java b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerS.java index ddb0094..feb8162 100644 --- a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerS.java +++ b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/po/influx/DataHarmpowerS.java @@ -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; diff --git a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/pojo/dto/PqsCommunicateDto.java b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/pojo/dto/PqsCommunicateDto.java index 6271716..4b7f81a 100644 --- a/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/pojo/dto/PqsCommunicateDto.java +++ b/data-processing/data-processing-api/src/main/java/com/njcn/dataProcess/pojo/dto/PqsCommunicateDto.java @@ -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; } diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/LnDataDealServiceImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/LnDataDealServiceImpl.java index 572771c..4e1726a 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/LnDataDealServiceImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/LnDataDealServiceImpl.java @@ -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> 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 lineIdList = lnDataDTO.getDataVList().stream().map(DataVDTO::getLineid).distinct().collect(Collectors.toList()); List data = lineFeignClient.getAllLine(lineIdList).getData(); + + Map map = data.stream().collect(Collectors.toMap(LineDeviceStateVO::getId, temp -> temp.getPids().split(",")[4])); lnDataDTO.getDataVList().forEach(temp->{ temp.setDevId(map.get(temp.getLineid())); }); Map> collect = lnDataDTO.getDataVList().stream().collect(Collectors.groupingBy(DataVDTO::getDevId)); + List> 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(); } diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlickerImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlickerImpl.java index e9e804e..c4a22a3 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlickerImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlickerImpl.java @@ -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 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 PHASE_MAPPING = new HashMap() {{ put("AB", "A"); put("BC", "B"); put("CA", "C"); put("M", "T"); }}; + @Resource + private RedisUtil redisUtil; + private static final Set LINE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T"))); + private static final Set PHASE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T"))); @Override @@ -229,6 +238,13 @@ public class InfluxdbDataFlickerImpl extends MppServiceImpl getMinuteData(LineCountEvaluateParam lineParam) { List result = new ArrayList<>(); + List data = new ArrayList<>(); + //获取监测点、接线方式数据 + Type type = new TypeToken>(){}.getType(); + Map 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 { - String newType = PHASE_MAPPING.get(item.getPhasicType()); - if (newType != null) { - item.setPhasicType(newType); - } - }); + if (!Objects.isNull(map)) { + //现根据监测点分组,然后根据接线方式排除多于数据,在修改相别 + Map> 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 validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES; + List 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; } } diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlucImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlucImpl.java index 4b7ac26..486ace8 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlucImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataFlucImpl.java @@ -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 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 PHASE_MAPPING = new HashMap() {{ put("AB", "A"); put("BC", "B"); put("CA", "C"); put("M", "T"); }}; + @Resource + private RedisUtil redisUtil; + private static final Set LINE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T"))); + private static final Set PHASE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T"))); @Override public void batchInsertion(List dataFlucDTOList) { @@ -147,6 +156,13 @@ public class InfluxdbDataFlucImpl extends MppServiceImpl getMinuteData(List lineList, String startTime, String endTime, Map> timeMap, Boolean dataType) { List dataList; List result = new ArrayList<>(); + List data = new ArrayList<>(); + //获取监测点、接线方式数据 + Type type = new TypeToken>(){}.getType(); + Map 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 { - String newType = PHASE_MAPPING.get(item.getPhasicType()); - if (newType != null) { - item.setPhasicType(newType); - } - }); + if (!Objects.isNull(map)) { + //现根据监测点分组,然后根据接线方式排除多于数据,在修改相别 + Map> 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 validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES; + List 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; } } diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmRateVImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmRateVImpl.java index 052e359..615ac8a 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmRateVImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmRateVImpl.java @@ -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 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 PHASE_MAPPING = new HashMap() {{ put("AB", "A"); put("BC", "B"); put("CA", "C"); put("M", "T"); }}; + @Resource + private RedisUtil redisUtil; + private static final Set LINE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T"))); + private static final Set PHASE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T"))); @Override public List getRawData(LineCountEvaluateParam lineParam) { @@ -244,6 +253,13 @@ public class InfluxdbDataHarmRateVImpl extends MppServiceImpl getMinuteData(LineCountEvaluateParam lineParam) { List dataList; List result = new ArrayList<>(); + List data = new ArrayList<>(); + //获取监测点、接线方式数据 + Type type = new TypeToken>(){}.getType(); + Map 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 { - String newType = PHASE_MAPPING.get(item.getPhasicType()); - if (newType != null) { - item.setPhasicType(newType); - } - }); + if (!Objects.isNull(map)) { + //现根据监测点分组,然后根据接线方式排除多于数据,在修改相别 + Map> 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 validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES; + List 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; } } diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmphasicVImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmphasicVImpl.java index 8e297c5..719849e 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmphasicVImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmphasicVImpl.java @@ -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 PHASE_MAPPING = new HashMap() {{ put("AB", "A"); put("BC", "B"); put("CA", "C"); put("M", "T"); }}; + @Resource + private RedisUtil redisUtil; + private static final Set LINE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T"))); + private static final Set PHASE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T"))); @Override @@ -205,6 +215,13 @@ public class InfluxdbDataHarmphasicVImpl extends MppServiceImpl getMinuteData(List lineList, String startTime, String endTime, Map> timeMap, Boolean dataType) { List dataList; List result = new ArrayList<>(); + List data = new ArrayList<>(); + //获取监测点、接线方式数据 + Type type = new TypeToken>(){}.getType(); + Map 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 { - String newType = PHASE_MAPPING.get(item.getPhasicType()); - if (newType != null) { - item.setPhasicType(newType); - } - }); + if (!Objects.isNull(map)) { + //现根据监测点分组,然后根据接线方式排除多于数据,在修改相别 + Map> 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 validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES; + List 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; } } diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerPImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerPImpl.java index 9b567a4..989f9a8 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerPImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerPImpl.java @@ -192,6 +192,7 @@ public class InfluxdbDataHarmpowerPImpl extends MppServiceImpl 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{ CommonMinuteDto.ValueType value = new CommonMinuteDto.ValueType(); value.setValueType(valueType); - List> lists; if (Objects.equals(phasicType, "T") && Objects.equals(lineParam.getType(), 2)) { lists = extractDataLists(valueTypeList, "Tot"); @@ -188,11 +187,13 @@ public class InfluxdbDataHarmpowerQImpl extends MppServiceImpl 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) diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerSImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerSImpl.java index 11df88a..80866ea 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerSImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataHarmpowerSImpl.java @@ -187,11 +187,13 @@ public class InfluxdbDataHarmpowerSImpl extends MppServiceImpl 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) diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataPltImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataPltImpl.java index 960dbcc..95cdbdf 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataPltImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataPltImpl.java @@ -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 PHASE_MAPPING = new HashMap() {{ put("AB", "A"); put("BC", "B"); put("CA", "C"); put("M", "T"); }}; + @Resource + private RedisUtil redisUtil; + private static final Set LINE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T"))); + private static final Set PHASE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T"))); @Override public void batchInsertion(List dataPltDTOList) { @@ -152,6 +162,13 @@ public class InfluxdbDataPltImpl extends MppServiceImpl getMinuteDataPlt(LineCountEvaluateParam lineParam) { List dataList; List result = new ArrayList<>(); + List data = new ArrayList<>(); + //获取监测点、接线方式数据 + Type type = new TypeToken>(){}.getType(); + Map 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 { - String newType = PHASE_MAPPING.get(item.getPhasicType()); - if (newType != null) { - item.setPhasicType(newType); - } - }); + if (!Objects.isNull(map)) { + //现根据监测点分组,然后根据接线方式排除多于数据,在修改相别 + Map> 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 validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES; + List 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; } } diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataVImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataVImpl.java index 89e91e6..20d47d4 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataVImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbDataVImpl.java @@ -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 implements IDataV { private final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault()); - private static final Map PHASE_MAPPING = new HashMap() {{ put("AB", "A"); put("BC", "B"); put("CA", "C"); put("M", "T"); }}; - @Resource private DataVMapper dataVMapper; + @Resource + private RedisUtil redisUtil; + private static final Set LINE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("AB", "BC", "CA", "T"))); + private static final Set PHASE_VOLTAGE_TYPES = + Collections.unmodifiableSet(new HashSet<>(Arrays.asList("A", "B", "C", "T"))); /** * 注意:influxdb不推荐采用in函数的方式批量查询监测点的数据,效率很低,容易造成崩溃,故每次单测点查询 @@ -437,6 +446,13 @@ public class InfluxdbDataVImpl extends MppServiceImpl getMinuteDataV(LineCountEvaluateParam lineParam) { List result = new ArrayList<>(); + List data = new ArrayList<>(); + //获取监测点、接线方式数据 + Type type = new TypeToken>(){}.getType(); + Map 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 { - String newType = PHASE_MAPPING.get(item.getPhasicType()); - if (newType != null) { - item.setPhasicType(newType); - } - }); + if (!Objects.isNull(map)) { + //现根据监测点分组,然后根据接线方式排除多于数据,在修改相别 + Map> lineMap = result.stream().collect(Collectors.groupingBy(DataV::getLineId)); + lineMap.forEach((k,v)->{ + if (Objects.isNull(map.get(k))) { + return; + } + //这边需要特殊处理下,将线电压数据赋值 + Map 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 validPhasicTypes = (conType != 0) ? LINE_VOLTAGE_TYPES : PHASE_VOLTAGE_TYPES; + List 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 result, InfluxQueryWrapper influxQueryWrapper, LineCountEvaluateParam lineParam) { diff --git a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbPqsCommunicateImpl.java b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbPqsCommunicateImpl.java index d56058e..497a3f3 100644 --- a/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbPqsCommunicateImpl.java +++ b/data-processing/data-processing-boot/src/main/java/com/njcn/dataProcess/service/impl/influxdb/InfluxdbPqsCommunicateImpl.java @@ -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 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 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); + + } diff --git a/message/message-api/src/main/java/com/njcn/message/message/DeviceRebootMessage.java b/message/message-api/src/main/java/com/njcn/message/message/DeviceRebootMessage.java index 4d24027..20fd7a3 100644 --- a/message/message-api/src/main/java/com/njcn/message/message/DeviceRebootMessage.java +++ b/message/message-api/src/main/java/com/njcn/message/message/DeviceRebootMessage.java @@ -31,6 +31,8 @@ public class DeviceRebootMessage { private String ip; //终端型号 private String devType; + //前置类型 + private String comType; //挂载单位 private String org_name; //组织名称 diff --git a/message/message-boot/pom.xml b/message/message-boot/pom.xml index a5e3fc3..0e5f597 100644 --- a/message/message-boot/pom.xml +++ b/message/message-boot/pom.xml @@ -100,12 +100,34 @@ spring-boot-starter-websocket 2.7.12 + + com.njcn + mq-spring-boot-starter + 1.0.0-SNAPSHOT + + + org.springframework.boot + spring-boot-starter-test + 2.7.12 + test + + + junit + junit + 4.13.2 + test + messageboot + + org.apache.maven.plugins + maven-surefire-plugin + 2.22.2 + org.apache.maven.plugins maven-compiler-plugin diff --git a/message/message-boot/src/main/java/com/njcn/message/consumer/FrontDataConsumer.java b/message/message-boot/src/main/java/com/njcn/message/consumer/FrontDataConsumer.java index 335b9b1..0b12afd 100644 --- a/message/message-boot/src/main/java/com/njcn/message/consumer/FrontDataConsumer.java +++ b/message/message-boot/src/main/java/com/njcn/message/consumer/FrontDataConsumer.java @@ -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", diff --git a/message/message-boot/src/main/java/com/njcn/message/mq/FrontDataMqListener.java b/message/message-boot/src/main/java/com/njcn/message/mq/FrontDataMqListener.java new file mode 100644 index 0000000..2219940 --- /dev/null +++ b/message/message-boot/src/main/java/com/njcn/message/mq/FrontDataMqListener.java @@ -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)); + } +} diff --git a/message/message-boot/src/main/java/com/njcn/message/mq/RecallMqSender.java b/message/message-boot/src/main/java/com/njcn/message/mq/RecallMqSender.java new file mode 100644 index 0000000..22998e2 --- /dev/null +++ b/message/message-boot/src/main/java/com/njcn/message/mq/RecallMqSender.java @@ -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); + } +} diff --git a/message/message-boot/src/test/java/com/njcn/message/mq/MqSliceWiringTest.java b/message/message-boot/src/test/java/com/njcn/message/mq/MqSliceWiringTest.java new file mode 100644 index 0000000..de2d75f --- /dev/null +++ b/message/message-boot/src/test/java/com/njcn/message/mq/MqSliceWiringTest.java @@ -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:6379(docker redis-stream-it);仅上/下行收发用例需要, + * 纯装配用例不需要(Lettuce 懒连接,不触发实际连接)。 + * + *

装配一律用 {@code withConfiguration(AutoConfigurations.of(...))},让 Spring 按 + * 真实 autoconfig 顺序(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() { + // 纯装配检查,不依赖真实 Redis(Lettuce 懒连接)。 + 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;不装 RedisAutoConfiguration,redis 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 fields = new HashMap(); + 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_Topic(MqListenerRegistry.start() 已在 context 启动时订阅该 stream) + redis.opsForStream().add( + StreamRecords.newRecord().in("LN_Topic").ofMap(fields) + ); + + // 等待消费线程处理(blockMs 默认 2000ms,timeout=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> records = + redis.opsForStream().read( + StreamReadOptions.empty().count(1L), + StreamOffset.create("n1_recall_Topic", ReadOffset.from("0-0")) + ); + assertThat(records).isNotEmpty(); + + Map 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"); + }); + } +}