Kaynağa Gözat

运行部分修正

lisonglin 2 hafta önce
ebeveyn
işleme
35e680b929

+ 113 - 6
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/controller/optimization/IndicatorController.java

@@ -1,6 +1,10 @@
 package com.hdkj.lt.bf.controller.optimization;
 
+import com.alibaba.fastjson.JSON;
+import com.hdkj.fhzg.optimization.request.MaintOrgRequest;
 import com.hdkj.hussar.ApiResponse;
+import com.hdkj.lt.bf.adapter.SiStateEstimationClient;
+import com.hdkj.lt.bf.common.SeSnapshotConstants;
 import com.hdkj.lt.bf.entity.vo.IndicatorDashboardVO;
 import com.hdkj.lt.bf.entity.vo.PageResult;
 import com.hdkj.lt.bf.entity.vo.SeCapacityDashboardVO;
@@ -12,14 +16,19 @@ import com.hdkj.lt.bf.entity.vo.SeVoltageDashboardVO;
 import com.hdkj.lt.bf.service.IndicatorDashboardService;
 import com.hdkj.lt.bf.service.SeAlarmService;
 import com.hdkj.lt.bf.service.SeIndicatorService;
+import com.hdkj.lt.bf.service.SeSnapshotService;
+import com.hdkj.lt.core.bizms.modle.dto.StateEstimation;
 import com.hdkj.lt.core.mvc.BaseController;
 import lombok.AllArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.StringUtils;
 import org.springframework.web.bind.annotation.PostMapping;
 import org.springframework.web.bind.annotation.RequestBody;
 import org.springframework.web.bind.annotation.RequestMapping;
 import org.springframework.web.bind.annotation.RestController;
 
+import java.util.ArrayList;
+import java.util.List;
 import java.util.Map;
 
 /**
@@ -39,6 +48,8 @@ public class IndicatorController extends BaseController {
     private final IndicatorDashboardService indicatorDashboardService;
     private final SeIndicatorService seIndicatorService;
     private final SeAlarmService seAlarmService;
+    private final SiStateEstimationClient siStateEstimationClient;
+    private final SeSnapshotService seSnapshotService;
 
     /**
      * 查询运行态仪表盘
@@ -106,14 +117,15 @@ public class IndicatorController extends BaseController {
     /**
      * 异常告警-线路重过载列表
      *
-     * @param params {id, type(2/3/4), feederName(可选), date(可选), orderBy(asc/desc,可选)}
+     * @param params {id, type(2/3/4), feederName(可选), startDate(可选), endDate(可选), orderBy(asc/desc,可选)}
      */
     @PostMapping("/alarm/overload")
     public ApiResponse<PageResult<SeAlarmListVO>> alarmOverload(@RequestBody Map<String, String> params) {
         String id = params.get("id");
         String type = params.get("type");
         String feederName = params.get("feederName");
-        String date = params.get("date");
+        String startDate = params.get("startDate");
+        String endDate = params.get("endDate");
         String orderBy = params.get("orderBy");
         String pageStr = params.get("page");
         String pageSizeStr = params.get("pageSize");
@@ -122,20 +134,21 @@ public class IndicatorController extends BaseController {
         }
         Integer page = pageStr != null ? Integer.parseInt(pageStr) : 1;
         Integer pageSize = pageSizeStr != null ? Integer.parseInt(pageSizeStr) : 20;
-        return ApiResponse.success(seAlarmService.queryOverloadAlarmList(id, Integer.parseInt(type), feederName, date, orderBy, page, pageSize));
+        return ApiResponse.success(seAlarmService.queryOverloadAlarmList(id, Integer.parseInt(type), feederName, startDate, endDate, orderBy, page, pageSize));
     }
 
     /**
      * 异常告警-电压越限列表
      *
-     * @param params {id, type(2/3/4), feederName(可选), date(可选), orderBy(asc/desc,可选)}
+     * @param params {id, type(2/3/4), feederName(可选), startDate(可选), endDate(可选), orderBy(asc/desc,可选)}
      */
     @PostMapping("/alarm/voltage")
     public ApiResponse<PageResult<SeAlarmListVO>> alarmVoltage(@RequestBody Map<String, String> params) {
         String id = params.get("id");
         String type = params.get("type");
         String feederName = params.get("feederName");
-        String date = params.get("date");
+        String startDate = params.get("startDate");
+        String endDate = params.get("endDate");
         String orderBy = params.get("orderBy");
         String pageStr = params.get("page");
         String pageSizeStr = params.get("pageSize");
@@ -144,7 +157,7 @@ public class IndicatorController extends BaseController {
         }
         Integer page = pageStr != null ? Integer.parseInt(pageStr) : 1;
         Integer pageSize = pageSizeStr != null ? Integer.parseInt(pageSizeStr) : 20;
-        return ApiResponse.success(seAlarmService.queryVoltageAlarmList(id, Integer.parseInt(type), feederName, date, orderBy, page, pageSize));
+        return ApiResponse.success(seAlarmService.queryVoltageAlarmList(id, Integer.parseInt(type), feederName, startDate, endDate, orderBy, page, pageSize));
     }
 
     /**
@@ -175,4 +188,98 @@ public class IndicatorController extends BaseController {
         }
         return ApiResponse.success(seAlarmService.queryFeederList(id, Integer.parseInt(type)));
     }
+
+    /**
+     * 手动触发状估断面拉取
+     * <p>
+     * 按县域调 SI queryEquipByCounty 拉取最新状估数据,聚合后送入 processSnapshot 解析。
+     * 传 snapTime 时启用覆盖更新:先删该断面时刻存量明细再重新写入。
+     *
+     * @param params {snapTime(可选, yyyy-MM-dd HH:mm:ss),传值时覆盖更新}
+     */
+    @PostMapping("/snapshot/manual")
+    public ApiResponse<String> manualSnapshot(@RequestBody Map<String, String> params) {
+        String snapTimeStr = params.get("snapTime");
+        boolean overwrite = StringUtils.isNotBlank(snapTimeStr);
+
+        log.info("手动状估拉取开始 overwrite={}", overwrite);
+        int successCount = 0;
+        List<String> errors = new ArrayList<>();
+
+        for (Map.Entry<String, String[]> county : SeSnapshotConstants.COUNTIES) {
+            String maintOrg = county.getKey();
+            String countyName = county.getValue()[1];
+            try {
+                MaintOrgRequest request = new MaintOrgRequest();
+                request.setMaintOrg(maintOrg);
+
+                ApiResponse<List<StateEstimation>> resp = siStateEstimationClient.queryEquipByCounty(request);
+                if (resp == null || !resp.isSuccess() || resp.getData() == null || resp.getData().isEmpty()) {
+                    errors.add(countyName + "无数据");
+                    log.warn("手动状估拉取 县域{} 无数据, resp={}", countyName, resp != null ? resp.isSuccess() : "null");
+                    continue;
+                }
+
+                StateEstimation merged = mergeCountyResults(resp.getData());
+                if (merged == null) {
+                    errors.add(countyName + "聚合后无数据");
+                    continue;
+                }
+
+                String json = JSON.toJSONString(merged);
+                if (overwrite) {
+                    seSnapshotService.processSnapshotOverwrite(json);
+                } else {
+                    seSnapshotService.processSnapshot(json);
+                }
+                successCount++;
+                log.info("手动状估拉取完成 county={}, pointTime={}", countyName, merged.getPointTime());
+            } catch (Exception e) {
+                log.error("手动状估拉取异常 county={}", countyName, e);
+                errors.add(countyName + "异常:" + e.getMessage());
+            }
+        }
+
+        StringBuilder msg = new StringBuilder();
+        msg.append("成功处理").append(successCount).append("个县域");
+        if (!errors.isEmpty()) {
+            msg.append(",").append(errors.size()).append("个异常:").append(String.join("; ", errors));
+        }
+        return ApiResponse.success(msg.toString());
+    }
+
+    private StateEstimation mergeCountyResults(List<StateEstimation> list) {
+        if (list == null || list.isEmpty()) return null;
+
+        StateEstimation merged = new StateEstimation();
+        List<StateEstimation.StateEstimationCommonResult> allFeeders = new ArrayList<>();
+        List<StateEstimation.StateEstimationTransResult> allMvtrans = new ArrayList<>();
+        List<StateEstimation.StateEstimationSwitchResult> allSwitches = new ArrayList<>();
+        List<StateEstimation.StateEstimationCommonResult> allEquips = new ArrayList<>();
+        List<StateEstimation.StateEstimationSegmentResult> allSegments = new ArrayList<>();
+        List<String> pointTime = null;
+
+        for (StateEstimation se : list) {
+            if (se == null) continue;
+            if (pointTime == null && se.getPointTime() != null) pointTime = se.getPointTime();
+            if (se.getPeriodFeederSeResult() != null) allFeeders.addAll(se.getPeriodFeederSeResult());
+            if (se.getPeriodMVTransSeResult() != null) allMvtrans.addAll(se.getPeriodMVTransSeResult());
+            if (se.getPeriodSwitchSeResult() != null) allSwitches.addAll(se.getPeriodSwitchSeResult());
+            if (se.getPeriodEquipResult() != null) allEquips.addAll(se.getPeriodEquipResult());
+            if (se.getPeriodSegmentSeResult() != null) allSegments.addAll(se.getPeriodSegmentSeResult());
+        }
+
+        if (allFeeders.isEmpty() && allMvtrans.isEmpty()) {
+            log.warn("手动状估拉取 聚合后无任何数据");
+            return null;
+        }
+
+        merged.setPeriodFeederSeResult(allFeeders);
+        merged.setPeriodMVTransSeResult(allMvtrans);
+        merged.setPeriodSwitchSeResult(allSwitches);
+        merged.setPeriodEquipResult(allEquips);
+        merged.setPeriodSegmentSeResult(allSegments);
+        merged.setPointTime(pointTime);
+        return merged;
+    }
 }

+ 3 - 9
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/scheduler/SeSnapshotScheduler.java

@@ -15,16 +15,13 @@ import org.springframework.stereotype.Component;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.Map;
-import java.util.concurrent.CompletableFuture;
 
 /**
  * 状估断面定时拉取任务
  * <p>
- * 每15分钟调用同事的 queryEquipByCounty 接口,按县域拉取最新状估数据,
+ * 每15分钟调用同事的 queryEquipByCounty 接口,按县域串行拉取最新状估数据,
  * 聚合后喂给 SeSnapshotService.processSnapshot() 解析入库,
  * 完成状估数据→断面明细→状态机监测→告警事件的完整链路。
- * <p>
- * 县域间通过 CompletableFuture 并行拉取,互不阻塞。
  *
  * @author lsl
  * @since 2026-07-26
@@ -39,17 +36,14 @@ public class SeSnapshotScheduler {
 
     @Scheduled(cron = "0 */15 * * * ?")
     public void pullStateEstimation() {
-        log.info("【状估拉取】开始定时拉取状估数据");
+        log.info("【状估拉取】开始串行拉取状估数据");
 
-        List<CompletableFuture<Void>> futures = new ArrayList<>();
         for (Map.Entry<String, String[]> county : SeSnapshotConstants.COUNTIES) {
             String maintOrg = county.getKey();
             String countySid = county.getValue()[0];
             String countyName = county.getValue()[1];
-            futures.add(CompletableFuture.runAsync(() ->
-                    pullByCounty(maintOrg, countySid, countyName)));
+            pullByCounty(maintOrg, countySid, countyName);
         }
-        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
     }
 
     private void pullByCounty(String maintOrg, String countySid, String countyName) {

+ 8 - 6
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/SeAlarmService.java

@@ -22,23 +22,25 @@ import java.util.List;
 public interface SeAlarmService {
 
     /**
-     * 分页查询指定日期重过载告警列表
+     * 分页查询指定日期范围重过载告警列表
      *
-     * @param date  查询日期 yyyy-MM-dd,为空则当天
+     * @param startDate 查询起始日期 yyyy-MM-dd,为空则当天
+     * @param endDate   查询截止日期 yyyy-MM-dd,为空则同 startDate
      * @param orderDir 排序方向 asc/desc,为空则 desc
      */
     PageResult<SeAlarmListVO> queryOverloadAlarmList(String id, Integer type, String feederName,
-                                                      String date, String orderDir,
+                                                      String startDate, String endDate, String orderDir,
                                                       Integer page, Integer pageSize);
 
     /**
-     * 分页查询指定日期电压越限告警列表
+     * 分页查询指定日期范围电压越限告警列表
      *
-     * @param date  查询日期 yyyy-MM-dd,为空则当天
+     * @param startDate 查询起始日期 yyyy-MM-dd,为空则当天
+     * @param endDate   查询截止日期 yyyy-MM-dd,为空则同 startDate
      * @param orderDir 排序方向 asc/desc,为空则 desc
      */
     PageResult<SeAlarmListVO> queryVoltageAlarmList(String id, Integer type, String feederName,
-                                                     String date, String orderDir,
+                                                     String startDate, String endDate, String orderDir,
                                                      Integer page, Integer pageSize);
 
     /**

+ 10 - 1
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/SeSnapshotService.java

@@ -16,12 +16,21 @@ package com.hdkj.lt.bf.service;
 public interface SeSnapshotService {
 
     /**
-     * 处理一个状估断面(单县 JSONObject)
+     * 处理一个状估断面(单县 JSONObject)<br>
+     * 正常模式:已有数据则跳过(防重)。
      *
      * @param snapshotJson 状估JSON字符串(单县)
      */
     void processSnapshot(String snapshotJson);
 
+    /**
+     * 处理一个状估断面(单县 JSONObject)<br>
+     * 覆盖模式:先删该断面时刻的存量明细,再重新写入并触发状态机。
+     *
+     * @param snapshotJson 状估JSON字符串(单县)
+     */
+    void processSnapshotOverwrite(String snapshotJson);
+
     /**
      * 处理一批状估断面(多县 JSONArray)
      *

+ 16 - 12
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/SeAlarmServiceImpl.java

@@ -56,20 +56,22 @@ public class SeAlarmServiceImpl implements SeAlarmService {
 
     @Override
     public PageResult<SeAlarmListVO> queryOverloadAlarmList(String id, Integer type, String feederName,
-                                                             String date, String orderDir,
+                                                             String startDate, String endDate, String orderDir,
                                                              Integer page, Integer pageSize) {
         if (page == null || page < 1) page = 1;
         if (pageSize == null || pageSize < 1) pageSize = 20;
 
-        LocalDate queryDate = parseDate(date);
-        LocalDateTime dayStart = queryDate.atStartOfDay();
-        LocalDateTime dayEnd = queryDate.plusDays(1).atStartOfDay();
+        LocalDate start = parseDate(startDate);
+        LocalDate end = parseDate(endDate);
+        if (end == null || end.isBefore(start)) end = start;
+        LocalDateTime rangeStart = start.atStartOfDay();
+        LocalDateTime rangeEnd = end.plusDays(1).atStartOfDay();
 
         Page<FhzgSeCurrentEvent> mpPage = new Page<>(page, pageSize);
 
         LambdaQueryWrapper<FhzgSeCurrentEvent> wrapper = new LambdaQueryWrapper<FhzgSeCurrentEvent>()
-                .ge(FhzgSeCurrentEvent::getFirstOverTime, dayStart)
-                .lt(FhzgSeCurrentEvent::getFirstOverTime, dayEnd)
+                .ge(FhzgSeCurrentEvent::getFirstOverTime, rangeStart)
+                .lt(FhzgSeCurrentEvent::getFirstOverTime, rangeEnd)
                 .eq(FhzgSeCurrentEvent::getStatus, 1)
                 .in(FhzgSeCurrentEvent::getAlarmType, OVERLOAD_TYPES);
 
@@ -105,20 +107,22 @@ public class SeAlarmServiceImpl implements SeAlarmService {
 
     @Override
     public PageResult<SeAlarmListVO> queryVoltageAlarmList(String id, Integer type, String feederName,
-                                                            String date, String orderDir,
+                                                            String startDate, String endDate, String orderDir,
                                                             Integer page, Integer pageSize) {
         if (page == null || page < 1) page = 1;
         if (pageSize == null || pageSize < 1) pageSize = 20;
 
-        LocalDate queryDate = parseDate(date);
-        LocalDateTime dayStart = queryDate.atStartOfDay();
-        LocalDateTime dayEnd = queryDate.plusDays(1).atStartOfDay();
+        LocalDate start = parseDate(startDate);
+        LocalDate end = parseDate(endDate);
+        if (end == null || end.isBefore(start)) end = start;
+        LocalDateTime rangeStart = start.atStartOfDay();
+        LocalDateTime rangeEnd = end.plusDays(1).atStartOfDay();
 
         Page<FhzgSeVoltageEvent> mpPage = new Page<>(page, pageSize);
 
         LambdaQueryWrapper<FhzgSeVoltageEvent> wrapper = new LambdaQueryWrapper<FhzgSeVoltageEvent>()
-                .ge(FhzgSeVoltageEvent::getFirstOverTime, dayStart)
-                .lt(FhzgSeVoltageEvent::getFirstOverTime, dayEnd)
+                .ge(FhzgSeVoltageEvent::getFirstOverTime, rangeStart)
+                .lt(FhzgSeVoltageEvent::getFirstOverTime, rangeEnd)
                 .eq(FhzgSeVoltageEvent::getStatus, 1)
                 .in(FhzgSeVoltageEvent::getAlarmLevel, "feeder", "mvtrans");
 

+ 25 - 3
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/SeSnapshotServiceImpl.java

@@ -98,14 +98,12 @@ public class SeSnapshotServiceImpl implements SeSnapshotService {
         }
 
         if (!batch.isEmpty()) {
-            // 防重:过滤掉该断面时刻已存在的 device_id
+            // 防重:查该断面时刻已有的所有 device_id(精确到秒,无需再 IN 过滤)
             Set<String> existingIds = new HashSet<>();
             snapshotDetailMapper.selectList(
                     new LambdaQueryWrapper<FhzgSeSnapshotDetail>()
                             .select(FhzgSeSnapshotDetail::getDeviceId)
                             .eq(FhzgSeSnapshotDetail::getSnapshotTime, snapTime)
-                            .in(FhzgSeSnapshotDetail::getDeviceId,
-                                    batch.stream().map(FhzgSeSnapshotDetail::getDeviceId).collect(java.util.stream.Collectors.toList()))
             ).forEach(e -> existingIds.add(e.getDeviceId()));
 
             int inserted = 0;
@@ -507,4 +505,28 @@ public class SeSnapshotServiceImpl implements SeSnapshotService {
                 buildFeederMappingRecursive(node.getChildren(), map, subsId, countyId);
         }
     }
+
+    // ============================================================
+    // 覆盖更新
+    // ============================================================
+
+    @Override
+    @Transactional(rollbackFor = Exception.class)
+    public void processSnapshotOverwrite(String snapshotJson) {
+        JSONObject root = JSON.parseObject(snapshotJson);
+        if (root == null) return;
+
+        JSONArray pointTimes = root.getJSONArray("pointTime");
+        if (pointTimes == null || pointTimes.isEmpty()) return;
+        LocalDateTime snapTime = LocalDateTime.parse(pointTimes.getString(0), DT_FMT);
+
+        // 先删该断面时刻的存量明细(覆盖更新)
+        int deleted = snapshotDetailMapper.delete(
+                new LambdaQueryWrapper<FhzgSeSnapshotDetail>()
+                        .eq(FhzgSeSnapshotDetail::getSnapshotTime, snapTime));
+        log.info("状估覆盖更新 snapTime={}, 删除存量{}条", snapTime, deleted);
+
+        // 重新走正常流程写入(此时防重逻辑不会跳过任何记录)
+        processSnapshot(snapshotJson);
+    }
 }