改造:监测点数据完整性_日表

This commit is contained in:
2023-11-13 10:13:30 +08:00
parent 468b228d82
commit 856adce8c1
12 changed files with 222 additions and 127 deletions

View File

@@ -5,6 +5,7 @@ import com.baomidou.mybatisplus.annotation.TableName;
import com.github.jeffreyning.mybatisplus.anno.MppMultiId;
import com.njcn.db.bo.BaseEntity;
import java.io.Serializable;
import java.time.LocalDate;
import java.time.LocalDateTime;
import lombok.Data;
@@ -27,7 +28,7 @@ public class RStatIntegrityD {
private static final long serialVersionUID = 1L;
@MppMultiId
private LocalDateTime timeId;
private LocalDate timeId;
@MppMultiId
private String lineIndex;

View File

@@ -180,7 +180,7 @@ logging:
whitelist:
urls:
# - /**
- /**
- /user-boot/user/generateSm2Key
- /user-boot/theme/getTheme
- /user-boot/user/updateFirstPassword

View File

@@ -1,6 +1,7 @@
package com.njcn.influx.imapper;
import com.njcn.influx.base.InfluxDbBaseMapper;
import com.njcn.influx.pojo.bo.MeasurementCount;
import com.njcn.influx.pojo.po.DataV;
import com.njcn.influx.query.InfluxQueryWrapper;
@@ -18,4 +19,6 @@ public interface DataVMapper extends InfluxDbBaseMapper<DataV> {
List<DataV> getStatisticsByWraper(InfluxQueryWrapper influxQueryWrapper);
List<MeasurementCount> getMeasurementCount(InfluxQueryWrapper influxQueryWrapper);
}

View File

@@ -0,0 +1,33 @@
package com.njcn.influx.pojo.bo;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import com.njcn.common.utils.serializer.InstantDateSerializer;
import lombok.Data;
import org.influxdb.annotation.Column;
import org.influxdb.annotation.Measurement;
import java.io.Serializable;
import java.time.Instant;
/**
* 类的介绍:
*
* @author xuyang
* @version 1.0.0
* @createTime 2023/11/10 16:17
*/
@Data
@Measurement(name = "data_v")
public class MeasurementCount implements Serializable {
@Column(name = "time")
@JsonSerialize(using = InstantDateSerializer.class)
private Instant time;
@Column(name = "line_id")
private String lineId;
@Column(name = "freq")
private String freq;
}

View File

@@ -185,5 +185,10 @@ public interface InfluxDBTableConstant {
String AVG = "AVG";
String CP95 = "CP95";
/**
* 每天固定时间分钟
*/
Integer DAY_MINUTE = 1440;
}

View File

@@ -111,12 +111,23 @@
<artifactId>liteflow-spring-boot-starter</artifactId>
<version>2.11.2</version>
</dependency>
<dependency>
<groupId>com.njcn</groupId>
<artifactId>event-api</artifactId>
<version>1.0.0</version>
<scope>compile</scope>
</dependency>
<dependency>
<groupId>com.yomahub</groupId>
<artifactId>liteflow-rule-nacos</artifactId>
<version>2.11.2</version>
</dependency>
<!-- <dependency>-->
<!-- <groupId>com.yomahub</groupId>-->
<!-- <artifactId>liteflow-rule-nacos</artifactId>-->
<!-- <version>2.11.0</version>-->
<!-- </dependency>-->
</dependencies>
<build>

View File

@@ -24,17 +24,13 @@ import lombok.RequiredArgsConstructor;
public class MeasurementExecutor extends BaseExecutor {
private final RMpMonitorEvaluateDService rMpMonitorEvaluateDService;
private final RMpEventDetailService rMpEventDetailService;
private final RMpEventDetailDService rMpEventDetailDService;
private final DayDataService dayDataService;
private final RStatAbnormalDService rStatAbnormalDService;
private final ROperatingMonitorService rOperatingMonitorService;
private final ROperatingMonitorMService rOperatingMonitorMService;
private final IntegrityService integrityService;
private final RMpPassRateDService rMpPassRateDService;
@@ -81,20 +77,6 @@ public class MeasurementExecutor extends BaseExecutor {
}
}
/**
* 算法名: 3.4.1.1-----监测点报表_日表
*
* @author xuyang
* @date 2023年11月09日 10:08
*/
@LiteflowMethod(value = LiteFlowMethodEnum.IS_ACCESS, nodeId = "dataToDay", nodeType = NodeTypeEnum.COMMON)
public boolean dataToDayAccess(NodeComponent bindCmp) {
return isAccess(bindCmp);
}
@LiteflowMethod(value = LiteFlowMethodEnum.PROCESS, nodeId = "dataToDay", nodeType = NodeTypeEnum.COMMON)
public void dataToDayProcess(NodeComponent bindCmp) {
dayDataService.dataToDayHandler(bindCmp.getRequestData());
}
/**
* 3.3.1.2. 监测点数据异常_日表
* @param bindCmp
@@ -169,4 +151,49 @@ public class MeasurementExecutor extends BaseExecutor {
}
}
}
/********************************************算法负责人:xy***********************************************************/
/**
* 算法名: 3.4.1.1-----监测点报表_日表
*
* @author xuyang
* @date 2023年11月09日 10:08
*/
@LiteflowMethod(value = LiteFlowMethodEnum.IS_ACCESS, nodeId = "dataToDay", nodeType = NodeTypeEnum.COMMON)
public boolean dataToDayAccess(NodeComponent bindCmp) {
return isAccess(bindCmp);
}
@LiteflowMethod(value = LiteFlowMethodEnum.PROCESS, nodeId = "dataToDay", nodeType = NodeTypeEnum.COMMON)
public void dataToDayProcess(NodeComponent bindCmp) {
dayDataService.dataToDayHandler(bindCmp.getRequestData());
}
/**
* 算法名: 暂无-----监测点数据完整性_日表
*
* @author xuyang
* @date 2023年11月09日 10:08
*/
@LiteflowMethod(value = LiteFlowMethodEnum.IS_ACCESS, nodeId = "measurementIntegrity", nodeType = NodeTypeEnum.COMMON)
public boolean measurementIntegrityAccess(NodeComponent bindCmp) {
return isAccess(bindCmp);
}
@LiteflowMethod(value = LiteFlowMethodEnum.PROCESS, nodeId = "measurementIntegrity", nodeType = NodeTypeEnum.COMMON)
public void measurementIntegrityProcess(NodeComponent bindCmp) {
integrityService.dataIntegrity(bindCmp.getRequestData());
}
/********************************************算法负责人:xy结束***********************************************************/
}

View File

@@ -40,17 +40,6 @@ public class IntegrityController extends BaseController {
private final IntegrityService integrityService;
/* @Deprecated
@OperateInfo(info = LogEnum.BUSINESS_COMMON)
@PostMapping("/computeDataIntegrity")
@ApiOperation("数据完整性统计")
@ApiImplicitParam(name = "lineParam", value = "参数", required = true)
public HttpResult<String> computeDataIntegrity(@RequestBody @Validated LineParam lineParam){
String methodDescribe = getMethodDescribe("computeDataIntegrity");
String out = integrityService.computeDataIntegrity(lineParam);
return HttpResultUtil.assembleCommonResponseResult(CommonResponseEnum.SUCCESS, out, methodDescribe);
}*/
@OperateInfo(info = LogEnum.BUSINESS_COMMON)
@PostMapping("/dataIntegrity")
@ApiOperation("数据完整性统计(MySQL库)")
@@ -65,10 +54,10 @@ public class IntegrityController extends BaseController {
log.info(item+"-->开始执行");
startTime = item+" "+"00:00:00";
endTime = item+" "+"23:59:59";
integrityService.dataIntegrity(lineParam,startTime,endTime);
// integrityService.dataIntegrity(lineParam,startTime,endTime);
}
} else {
integrityService.dataIntegrity(lineParam,lineParam.getBeginTime(),lineParam.getEndTime());
// integrityService.dataIntegrity(lineParam,lineParam.getBeginTime(),lineParam.getEndTime());
}
return HttpResultUtil.assembleCommonResponseResult(CommonResponseEnum.SUCCESS, CommonResponseEnum.SUCCESS.getMessage(), methodDescribe);
}

View File

@@ -3,39 +3,29 @@ package com.njcn.prepare.harmonic.service.mysql.Impl.line;
import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.DateUtil;
import cn.hutool.core.date.LocalDateTimeUtil;
import com.alibaba.fastjson.JSON;
import com.alibaba.nacos.client.naming.utils.CollectionUtils;
import com.njcn.common.utils.HarmonicTimesUtil;
import com.njcn.harmonic.pojo.po.day.*;
import com.njcn.influx.constant.InfluxDbSqlConstant;
import com.njcn.influx.deprecated.InfluxDBPublicParam;
import com.njcn.influx.imapper.*;
import com.njcn.influx.imapper.day.*;
import com.njcn.influx.pojo.po.*;
import com.njcn.influx.pojo.po.day.*;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.influx.utils.InfluxDbUtils;
import com.njcn.prepare.bo.CalculatedParam;
import com.njcn.prepare.harmonic.service.mysql.day.*;
import com.njcn.prepare.harmonic.service.mysql.line.DayDataService;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import net.sf.cglib.core.Local;
import org.apache.commons.collections4.ListUtils;
import org.influxdb.dto.QueryResult;
import org.influxdb.impl.InfluxDBResultMapper;
import org.springframework.beans.BeanUtils;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
/**

View File

@@ -1,32 +1,33 @@
package com.njcn.prepare.harmonic.service.mysql.Impl.line;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.LocalDateTimeUtil;
import com.github.jeffreyning.mybatisplus.service.MppServiceImpl;
import com.njcn.common.pojo.enums.common.ServerEnum;
import com.njcn.device.biz.commApi.CommTerminalGeneralClient;
import com.njcn.device.biz.pojo.dto.LineDevGetDTO;
import com.njcn.device.biz.pojo.param.DeptGetLineParam;
import com.njcn.device.pq.api.LineFeignClient;
import com.njcn.device.pq.pojo.po.RStatIntegrityD;
import com.njcn.influx.deprecated.InfluxDBPublicParam;
import com.njcn.influx.constant.InfluxDbSqlConstant;
import com.njcn.influx.imapper.DataVMapper;
import com.njcn.influx.pojo.bo.MeasurementCount;
import com.njcn.influx.pojo.constant.InfluxDBTableConstant;
import com.njcn.influx.pojo.po.DataV;
import com.njcn.influx.query.InfluxQueryWrapper;
import com.njcn.influx.utils.InfluxDbUtils;
import com.njcn.prepare.bo.CalculatedParam;
import com.njcn.prepare.harmonic.mapper.mysql.day.RStatIntegrityDMapper;
import com.njcn.prepare.harmonic.pojo.param.LineParam;
import com.njcn.prepare.harmonic.service.mysql.line.IntegrityService;
import com.njcn.user.api.DeptFeignClient;
import com.njcn.user.pojo.po.Dept;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.influxdb.dto.QueryResult;
import org.influxdb.impl.InfluxDBResultMapper;
import org.apache.commons.collections4.ListUtils;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
import org.springframework.util.CollectionUtils;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.stream.Collectors;
/**
@@ -50,80 +51,109 @@ public class IntegrityServiceImpl extends MppServiceImpl<RStatIntegrityDMapper,
private final CommTerminalGeneralClient commTerminalGeneralClient;
/*@Override
@Async("asyncExecutor")
public String computeDataIntegrity(LineParam lineParam) {
List<LineDetail> lineDetailList;
if (CollUtil.isEmpty(lineParam.getLineIds())){
List<Overlimit> overLimitList = getAllLinesLimitData();
List<String> lineList = overLimitList.stream().map(Overlimit::getId).collect(Collectors.toList());
lineDetailList = lineFeignClient.getLineDetail(lineList).getData();
}else {
lineDetailList = lineFeignClient.getLineDetail(lineParam.getLineIds()).getData();
}
if (CollUtil.isEmpty(lineDetailList)){
return "未查询到监测点详情!";
}
Date dateOut = DateUtil.parse(lineParam.getDataDate());
List<String> records = new ArrayList<>();
for (LineDetail lineDetail :lineDetailList){
Map<String, String> tags = new HashMap<>();
Map<String, Object> fields = new HashMap<>();
tags.put("line_id",lineDetail.getId());
fields.put("due",DAY_MINUTE/lineDetail.getTimeInterval());
int dataCount = getDataCount(lineDetail.getId(),lineParam.getDataDate());
fields.put("real",dataCount);
Point point = influxDbUtils.pointBuilder("pqs_integrity", dateOut.getTime(), TimeUnit.MILLISECONDS,tags, fields);
BatchPoints batchPoints = BatchPoints.database(influxDbUtils.getDbName()).tag("line_id", lineDetail.getId()).retentionPolicy("").consistency(InfluxDB.ConsistencyLevel.ALL).build();
batchPoints.point(point);
records.add(batchPoints.lineProtocol());
}
//InfluxDb入表pqs_integrity
influxDbUtils.batchInsert(influxDbUtils.getDbName(),"", InfluxDB.ConsistencyLevel.ALL, records);
return "成功!";
}
*/
private final DataVMapper dataVMapper;
// @Override
// @Async("asyncExecutor")
// @Deprecated
// public void dataIntegrity(LineParam lineParam,String startTime,String endTime) {
// DateTimeFormatter df = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
// LocalDateTime dateTime = LocalDateTime.parse(startTime,df);
//
// List<LineDevGetDTO> lineDevGetDTOList = new ArrayList<>();
// if (CollUtil.isEmpty(lineParam.getLineIds())){
// Dept dept = deptFeignClient.getRootDept().getData();
//
// DeptGetLineParam deptGetLineParam = new DeptGetLineParam();
// deptGetLineParam.setDeptId(dept.getId());
// deptGetLineParam.setServerName(ServerEnum.HARMONIC.getName());
// List<String> monitorIds = commTerminalGeneralClient.getRunMonitorIds().getData();
// lineDevGetDTOList = commTerminalGeneralClient.getMonitorDetailList(monitorIds).getData();
// }else {
// lineDevGetDTOList = commTerminalGeneralClient.getMonitorDetailList(lineParam.getLineIds()).getData();
// }
// List<RStatIntegrityD> list = new ArrayList<>();
// for (LineDevGetDTO lineDetail :lineDevGetDTOList){
// int dataCount = getDataCount(lineDetail.getPointId(),startTime,endTime);
// RStatIntegrityD integrityDpo = new RStatIntegrityD();
// integrityDpo.setTimeId(dateTime);
// integrityDpo.setLineIndex(lineDetail.getPointId());
// integrityDpo.setDueTime(InfluxDBPublicParam.DAY_MINUTE/lineDetail.getInterval());
// integrityDpo.setRealTime(dataCount);
// list.add(integrityDpo);
// }
// this.saveOrUpdateBatchByMultiId(list,500);
// }
/********************************新算法************************************************/
@Override
@Async("asyncExecutor")
public void dataIntegrity(LineParam lineParam,String startTime,String endTime) {
DateTimeFormatter df = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
LocalDateTime dateTime = LocalDateTime.parse(startTime,df);
List<LineDevGetDTO> lineDevGetDTOList = new ArrayList<>();
if (CollUtil.isEmpty(lineParam.getLineIds())){
Dept dept = deptFeignClient.getRootDept().getData();
DeptGetLineParam deptGetLineParam = new DeptGetLineParam();
deptGetLineParam.setDeptId(dept.getId());
deptGetLineParam.setServerName(ServerEnum.HARMONIC.getName());
List<String> monitorIds = commTerminalGeneralClient.getRunMonitorIds().getData();
lineDevGetDTOList = commTerminalGeneralClient.getMonitorDetailList(monitorIds).getData();
}else {
lineDevGetDTOList = commTerminalGeneralClient.getMonitorDetailList(lineParam.getLineIds()).getData();
}
public void dataIntegrity(CalculatedParam calculatedParam) {
List<RStatIntegrityD> list = new ArrayList<>();
for (LineDevGetDTO lineDetail :lineDevGetDTOList){
int dataCount = getDataCount(lineDetail.getPointId(),startTime,endTime);
RStatIntegrityD integrityDpo = new RStatIntegrityD();
integrityDpo.setTimeId(dateTime);
integrityDpo.setLineIndex(lineDetail.getPointId());
integrityDpo.setDueTime(InfluxDBPublicParam.DAY_MINUTE/lineDetail.getInterval());
integrityDpo.setRealTime(dataCount);
list.add(integrityDpo);
List<String> lineIds = calculatedParam.getIdList();
String beginDay = LocalDateTimeUtil.format(
LocalDateTimeUtil.beginOfDay(LocalDateTimeUtil.parse(calculatedParam.getDataDate(), DatePattern.NORM_DATE_PATTERN)),
DatePattern.NORM_DATETIME_PATTERN
);
String endDay = LocalDateTimeUtil.format(
LocalDateTimeUtil.endOfDay(LocalDateTimeUtil.parse(calculatedParam.getDataDate(), DatePattern.NORM_DATE_PATTERN)),
DatePattern.NORM_DATETIME_PATTERN
);
//以尺寸100分片
List<List<String>> pendingIds = ListUtils.partition(lineIds,100);
for (List<String> pendingId : pendingIds) {
List<LineDevGetDTO> lineDevGetDTOList = commTerminalGeneralClient.getMonitorDetailList(pendingId).getData();
List<MeasurementCount> countList = this.getMeasurementCount(pendingId,beginDay,endDay);
list.addAll(
lineDevGetDTOList.stream()
.map(item -> {
RStatIntegrityD integrityDpo = new RStatIntegrityD();
integrityDpo.setTimeId(LocalDateTimeUtil.parseDate(calculatedParam.getDataDate(), DatePattern.NORM_DATE_PATTERN));
integrityDpo.setLineIndex(item.getPointId());
integrityDpo.setDueTime(InfluxDBTableConstant.DAY_MINUTE / item.getInterval());
integrityDpo.setRealTime(countList.stream()
.filter(item2 -> Objects.equals(item.getPointId(), item2.getLineId()))
.map(item2 -> (int) Double.parseDouble(item2.getFreq()))
.findFirst().orElse(0)
);
return integrityDpo;
})
.collect(Collectors.toList())
);
}
this.saveOrUpdateBatchByMultiId(list,500);
this.saveOrUpdateBatchByMultiId(list,1000);
}
private int getDataCount(String lineId,String startTime,String endTime){
QueryResult sqlResult = influxDbUtils.query("SELECT * FROM data_v WHERE time >= '" + startTime + "' and time <= '" + endTime + "' and line_id = '" + lineId + "' and phasic_type = 'T' and value_type = 'MAX' tz('Asia/Shanghai')");
InfluxDBResultMapper resultMapper = new InfluxDBResultMapper();
List<DataV> list = resultMapper.toPOJO(sqlResult, DataV.class);
if (CollectionUtils.isEmpty(list)){
return 0;
} else {
return list.size();
}
/**
* 获取data_v中各个监测点的数据总数
* @param lineIndex
* @param startTime
* @param endTime
* @return
*/
public List<MeasurementCount> getMeasurementCount(List<String> lineIndex, String startTime, String endTime) {
InfluxQueryWrapper influxQueryWrapper = new InfluxQueryWrapper(DataV.class,MeasurementCount.class);
influxQueryWrapper.regular(DataV::getLineId, lineIndex)
.eq(DataV::getValueType, InfluxDbSqlConstant.MAX)
.eq(DataV::getPhasicType, InfluxDBTableConstant.PHASE_TYPE_T)
.count(DataV::getFreq)
.groupBy(DataV::getLineId)
.between(DataV::getTime, startTime, endTime);
return dataVMapper.getMeasurementCount(influxQueryWrapper);
}
/********************************新算法结束************************************************/
// private int getDataCount(String lineId,String startTime,String endTime){
// QueryResult sqlResult = influxDbUtils.query("SELECT * FROM data_v WHERE time >= '" + startTime + "' and time <= '" + endTime + "' and line_id = '" + lineId + "' and phasic_type = 'T' and value_type = 'MAX' tz('Asia/Shanghai')");
// InfluxDBResultMapper resultMapper = new InfluxDBResultMapper();
// List<DataV> list = resultMapper.toPOJO(sqlResult, DataV.class);
// if (CollectionUtils.isEmpty(list)){
// return 0;
// } else {
// return list.size();
// }
// }
}

View File

@@ -18,6 +18,6 @@ public interface DayDataService {
* @date 2023/11/09 10:08
* @param calculatedParam 查询条件
*/
void dataToDayHandler(CalculatedParam calculatedParam);
void dataToDayHandler(CalculatedParam<?> calculatedParam);
}

View File

@@ -1,6 +1,6 @@
package com.njcn.prepare.harmonic.service.mysql.line;
import com.njcn.prepare.harmonic.pojo.param.LineParam;
import com.njcn.prepare.bo.CalculatedParam;
/**
* @author xiaoyao
@@ -9,7 +9,13 @@ import com.njcn.prepare.harmonic.pojo.param.LineParam;
*/
public interface IntegrityService {
//String computeDataIntegrity(LineParam lineParam);
// void dataIntegrity(LineParam lineParam,String startTime,String endTime);
void dataIntegrity(LineParam lineParam,String startTime,String endTime);
/***
* 监测点数据完整性_日表
* @author xuyang
* @date 2023/11/09 10:08
* @param calculatedParam 查询条件
*/
void dataIntegrity(CalculatedParam calculatedParam);
}