Просмотр исходного кода

配网电气模型接口多线程调用初始化

lisonglin 2 месяцев назад
Родитель
Сommit
0127ee890b

+ 2 - 2
.gitignore

@@ -7,7 +7,6 @@
 
 target/
 !.mvn/wrapper/maven-wrapper.jar
-
 ######################################################################
 # IDE
 
@@ -43,4 +42,5 @@ nbdist/
 
 !*/build/*.java
 !*/build/*.html
-!*/build/*.xml
+!*/build/*.xml
+**/src/test/*

+ 1 - 1
api/load-transfer-si-api/src/main/java/com/hdkj/lt/feign/ISimulationPlatformApiClient.java

@@ -33,6 +33,6 @@ public interface ISimulationPlatformApiClient {
 
     @ApiOperation(value = "获取配网电气模型")
     @PostMapping("/getElectricalModel")
-    ElectricalModelResponse getElectricalModel(@RequestBody ElectricalModelRequest request);
+    String getElectricalModel(@RequestBody ElectricalModelRequest request);
 
 }

+ 1 - 1
api/load-transfer-si-api/src/main/java/com/hdkj/lt/feign/fallback/SimulationPlatformApiClientFallBack.java

@@ -39,7 +39,7 @@ public class SimulationPlatformApiClientFallBack implements FallbackFactory<ISim
             }
 
             @Override
-            public ElectricalModelResponse getElectricalModel(ElectricalModelRequest request) {
+            public String getElectricalModel(ElectricalModelRequest request) {
                 log.error("配网电气模型接口-远程调用服务异常, 参数:{},异常:{}", request, cause.getMessage());
                 return null;
             }

+ 133 - 4
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/AutoTaskServiceImpl.java

@@ -16,8 +16,10 @@ import com.hdkj.lt.bf.entity.FhzgAutoTaskPlan;
 import com.hdkj.lt.bf.entity.FhzgAutoTaskPlanDetail;
 import com.hdkj.lt.bf.service.*;
 import com.hdkj.lt.core.bizms.modle.po.DwdShbDsFeederBase;
+import com.hdkj.lt.feign.ISimulationPlatformApiClient;
 import com.hdkj.lt.feign.IPlanCreateClient;
 import com.hdkj.lt.modle.dto.plan.AcceptPlanPowerCutListRequest;
+import com.hdkj.lt.modle.dto.simulation.ElectricalModelRequest;
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.commons.lang3.StringUtils;
@@ -31,6 +33,7 @@ import java.time.LocalDate;
 import java.time.LocalDateTime;
 import java.time.temporal.TemporalAdjusters;
 import java.util.*;
+import java.util.concurrent.*;
 import java.util.stream.Collectors;
 
 /**
@@ -44,14 +47,26 @@ import java.util.stream.Collectors;
 public class AutoTaskServiceImpl implements AutoTaskService {
 
     private final FhzgAutoTaskPlanService autoTaskPlanService;
-
     private final FhzgAutoTaskCountryInfoService autoTaskCountryInfoService;
-
     private final FhzgAutoTaskPlanDetailService autoTaskPlanDetailService;
-
     private final FhzgAutoTaskFeederChangeService autoTaskFeederChangeService;
-
     private final IPlanCreateClient planCreateClient;
+    private final ISimulationPlatformApiClient simulationPlatformApiClient;
+    private final FhzgElectricalModelSectionService electricalModelSectionService;
+
+    // ==================== 电气模型采集线程池 ====================
+    private static final ThreadPoolExecutor ELECTRICAL_MODEL_POOL = new ThreadPoolExecutor(
+            10, 20, 60, TimeUnit.SECONDS,
+            new LinkedBlockingQueue<>(200),
+            new ThreadFactory() {
+                private int count = 0;
+                @Override
+                public Thread newThread(Runnable r) {
+                    return new Thread(r, "elec-model-" + (++count));
+                }
+            },
+            new ThreadPoolExecutor.CallerRunsPolicy()
+    );
 
     @Override
     public List<AutoTaskFeederResponse> getCountryStatistic(AutoTaskQueryRequest request) {
@@ -468,4 +483,118 @@ public class AutoTaskServiceImpl implements AutoTaskService {
         }
     }
 
+    // ==================== 电气模型定时采集 ====================
+
+    /** 最大重试次数 */
+    private static final int MAX_RETRY = 3;
+    /** 失败率阈值:超过此比例不触发后续计算 */
+    private static final double FAIL_THRESHOLD = 1;
+
+    /**
+     * 定时采集电气模型数据(每15分钟)
+     * 入口方法,可由定时任务或手动调用触发
+     *
+     * @param cityCode   地市编码
+     * @param feederList 馈线ID列表
+     */
+    public void collectElectricalModel(String cityCode, List<String> feederList) {
+        if (feederList == null || feederList.isEmpty()) {
+            log.warn("电气模型采集-馈线列表为空, cityCode={}", cityCode);
+            return;
+        }
+
+        log.info("电气模型采集-开始, cityCode={}, 馈线数={}", cityCode, feederList.size());
+        long start = System.currentTimeMillis();
+
+        // 1. 并发调用推演平台,每条馈线一个CompletableFuture
+        List<CompletableFuture<String>> futures = new ArrayList<>();
+        for (String feederId : feederList) {
+            CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
+                return collectSingleFeederWithRetry(cityCode, feederId);
+            }, ELECTRICAL_MODEL_POOL);
+            futures.add(future);
+        }
+
+        // 2. 等待全部完成
+        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
+
+        // 3. 收集结果
+        int successCount = 0;
+        int failCount = 0;
+        List<String> failedFeederIds = new ArrayList<>();
+
+        for (int i = 0; i < futures.size(); i++) {
+            String feederId = feederList.get(i);
+            try {
+                String jsonStr = futures.get(i).get();
+                if (jsonStr != null) {
+                    // 4. 入库
+                    electricalModelSectionService.saveFromJson(jsonStr);
+                    successCount++;
+                } else {
+                    failCount++;
+                    failedFeederIds.add(feederId);
+                }
+            } catch (Exception e) {
+                failCount++;
+                failedFeederIds.add(feederId);
+                log.error("电气模型采集-入库失败, feederId={}", feederId, e);
+            }
+        }
+
+        long cost = System.currentTimeMillis() - start;
+        double failRate = feederList.isEmpty() ? 0 : (double) failCount / feederList.size();
+
+        log.info("电气模型采集-完成, cityCode={}, 成功={}, 失败={}, 耗时={}ms, 失败率={}",
+                cityCode, successCount, failCount, cost, String.format("%.2f%%", failRate * 100));
+
+        if (!failedFeederIds.isEmpty()) {
+            log.warn("电气模型采集-失败馈线: {}", failedFeederIds);
+        }
+
+        // 5. 失败率过高则不触发后续计算
+        if (failRate >= FAIL_THRESHOLD) {
+            log.error("电气模型采集-失败率超过{}%, 不触发指标计算, cityCode={}", (int)(FAIL_THRESHOLD * 100), cityCode);
+            return;
+        }
+
+        // 6. TODO: 发送MQ消息触发指标计算
+        // rocketMQTemplate.sendAsyncMessage(...);
+    }
+
+    /**
+     * 采集单条馈线(带重试)
+     * @return 原始JSON字符串,失败返回null
+     */
+    private String collectSingleFeederWithRetry(String cityCode, String feederId) {
+        for (int attempt = 1; attempt <= MAX_RETRY; attempt++) {
+            try {
+                ElectricalModelRequest request = new ElectricalModelRequest();
+                request.setCityId(cityCode);
+                request.setFeederList(Collections.singletonList(feederId));
+
+                String jsonStr = simulationPlatformApiClient.getElectricalModel(request);
+                if (jsonStr != null) {
+                    log.debug("电气模型采集-成功, feederId={}, attempt={}", feederId, attempt);
+                    return jsonStr;
+                }
+                log.warn("电气模型采集-返回null, feederId={}, attempt={}", feederId, attempt);
+            } catch (Exception e) {
+                log.warn("电气模型采集-调用异常, feederId={}, attempt={}/{}, error={}",
+                        feederId, attempt, MAX_RETRY, e.getMessage());
+            }
+            // 递增等待: 1s, 3s, 5s
+            if (attempt < MAX_RETRY) {
+                try {
+                    Thread.sleep(1000L * (2 * attempt - 1));
+                } catch (InterruptedException ie) {
+                    Thread.currentThread().interrupt();
+                    return null;
+                }
+            }
+        }
+        log.error("电气模型采集-重试耗尽, feederId={}", feederId);
+        return null;
+    }
+
 }

+ 8 - 72
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/FhzgElectricalModelSectionServiceImpl.java

@@ -53,69 +53,61 @@ public class FhzgElectricalModelSectionServiceImpl extends ServiceImpl<FhzgElect
             log.warn("配网电气模型JSON中缺少section_info");
             return;
         }
-
         FhzgElectricalModelSection section = parseSection(sectionJson);
         this.save(section);
-
         String sectionId = section.getSectionId();
         log.info("配网电气模型断面保存成功, sectionId={}", sectionId);
-
         // 2. 解析data_objects
         JSONObject dataObjects = root.getJSONObject("data_objects");
         if (dataObjects == null) {
             log.warn("配网电气模型JSON中缺少data_objects");
             return;
         }
-
         // 3. 保存节点
         List<FhzgElectricalModelNode> nodes = parseNodes(dataObjects.getJSONArray("nodes"), sectionId);
         if (CollUtil.isNotEmpty(nodes)) {
             nodeService.saveBatch(nodes, 500);
             log.info("保存节点 {} 条", nodes.size());
         }
-
         // 4. 保存导线段
         List<FhzgElectricalModelConductor> conductors = parseConductors(dataObjects.getJSONArray("conductors"), sectionId);
         if (CollUtil.isNotEmpty(conductors)) {
             conductorService.saveBatch(conductors, 500);
             log.info("保存导线段 {} 条", conductors.size());
         }
-
         // 5. 保存配电变压器
         List<FhzgElectricalModelTransformer> transformers = parseTransformers(dataObjects.getJSONArray("transformers"), sectionId);
         if (CollUtil.isNotEmpty(transformers)) {
             transformerService.saveBatch(transformers, 500);
             log.info("保存配电变压器 {} 条", transformers.size());
         }
-
         // 6. 保存开关设备
         List<FhzgElectricalModelSwitch> switches = parseSwitches(dataObjects.getJSONArray("switches"), sectionId);
         if (CollUtil.isNotEmpty(switches)) {
             switchService.saveBatch(switches, 500);
             log.info("保存开关设备 {} 条", switches.size());
         }
-
         // 7. 保存注入设备
         List<FhzgElectricalModelInjection> injections = parseInjections(dataObjects.getJSONArray("injection_devices"), sectionId);
         if (CollUtil.isNotEmpty(injections)) {
             injectionService.saveBatch(injections, 500);
             log.info("保存注入设备 {} 条", injections.size());
         }
-
-        log.info("配网电气模型数据保存完成, sectionId={}, 节点={}, 导线={}, 变压器={}, 开关={}, 注入={}",
-                sectionId,
-                CollUtil.size(nodes),
-                CollUtil.size(conductors),
-                CollUtil.size(transformers),
-                CollUtil.size(switches),
-                CollUtil.size(injections));
     }
 
+
+
+
+
+
+
+
     // ==================== 手动映射方法 ====================
 
     private FhzgElectricalModelSection parseSection(JSONObject json) {
         FhzgElectricalModelSection s = new FhzgElectricalModelSection();
         s.setSectionId(json.getString("id"));       // JSON id → entity sectionId
+        s.setFeederId(json.getString("feeder_id"));
         s.setName(json.getString("name"));
         s.setBaseKv(json.getBigDecimal("base_kv"));
         s.setBaseKva(json.getBigDecimal("base_kva"));
@@ -271,60 +263,4 @@ public class FhzgElectricalModelSectionServiceImpl extends ServiceImpl<FhzgElect
     private Integer boolToInt(Boolean val) {
         return val != null && val ? 1 : 0;
     }
-
-    // ==================== 本地测试入口 ====================
-
-    /**
-     * 本地测试入口:读取JSON文件,解析并打印统计信息
-     * 用法:直接运行main方法,修改JSON_PATH为你的文件路径
-     */
-    public static void main(String[] args) throws IOException {
-        // ========== 修改这里为你的JSON文件路径 ==========
-        String JSON_PATH = "C:/Users/lenovo/Documents/Downloads/standard_model_raw.json";
-        // ================================================
-
-        if (args.length > 0) {
-            JSON_PATH = args[0];
-        }
-
-        System.out.println("读取文件: " + JSON_PATH);
-        String jsonStr = new String(Files.readAllBytes(Paths.get(JSON_PATH)));
-        System.out.println("文件大小: " + jsonStr.length() + " 字符");
-
-        JSONObject root = JSON.parseObject(jsonStr);
-
-        // --- 使用与saveFromJson相同的解析逻辑 ---
-        FhzgElectricalModelSectionServiceImpl service = new FhzgElectricalModelSectionServiceImpl();
-
-        // 解析section_info
-        JSONObject sectionJson = root.getJSONObject("section_info");
-        FhzgElectricalModelSection section = service.parseSection(sectionJson);
-        System.out.println("\n=== 断面信息 ===");
-        System.out.println("sectionId: " + section.getSectionId());
-        System.out.println("name: " + section.getName());
-        System.out.println("baseKv: " + section.getBaseKv());
-        System.out.println("baseKva: " + section.getBaseKva());
-        System.out.println("time: " + section.getTime());
-
-        // 解析data_objects
-        JSONObject dataObjects = root.getJSONObject("data_objects");
-        String sectionId = section.getSectionId();
-
-        List<FhzgElectricalModelNode> nodes = service.parseNodes(dataObjects.getJSONArray("nodes"), sectionId);
-        List<FhzgElectricalModelConductor> conductors = service.parseConductors(dataObjects.getJSONArray("conductors"), sectionId);
-        List<FhzgElectricalModelTransformer> transformers = service.parseTransformers(dataObjects.getJSONArray("transformers"), sectionId);
-        List<FhzgElectricalModelSwitch> switches = service.parseSwitches(dataObjects.getJSONArray("switches"), sectionId);
-        List<FhzgElectricalModelInjection> injections = service.parseInjections(dataObjects.getJSONArray("injection_devices"), sectionId);
-
-        System.out.println("\n=== 数据统计 ===");
-        System.out.println("节点(nodes): " + nodes.size());
-        System.out.println("导线段(conductors): " + conductors.size());
-        System.out.println("配电变压器(transformers): " + transformers.size());
-        System.out.println("开关设备(switches): " + switches.size());
-        System.out.println("注入设备(injection_devices): " + injections.size());
-
-
-
-        System.out.println("\n解析完成,无异常。接入Spring后调用 saveFromJson(jsonStr) 即可入库。");
-    }
 }

+ 1 - 1
services/load-transfer-si/src/main/java/com/hdkj/lt/si/controller/simulation/SimulationController.java

@@ -45,7 +45,7 @@ public class SimulationController implements ISimulationPlatformApiClient {
     }
 
     @Override
-    public ElectricalModelResponse getElectricalModel(ElectricalModelRequest request) {
+    public String getElectricalModel(ElectricalModelRequest request) {
         return simulationService.getElectricalModel(request);
     }
 }

+ 2 - 2
services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/simulation/ISimulationService.java

@@ -22,7 +22,7 @@ public interface ISimulationService {
     /**
      * 获取配网电气模型
      * @param request 请求参数
-     * @return 电气模型响应
+     * @return 原始JSON字符串
      */
-    ElectricalModelResponse getElectricalModel(ElectricalModelRequest request);
+    String getElectricalModel(ElectricalModelRequest request);
 }

+ 5 - 5
services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/simulation/impl/SimulationServiceImpl.java

@@ -216,9 +216,9 @@ public class SimulationServiceImpl implements ISimulationService {
     }
 
     @Override
-    public ElectricalModelResponse getElectricalModel(ElectricalModelRequest request) {
+    public String getElectricalModel(ElectricalModelRequest request) {
         // TODO: 配置化接口地址
-        String url =  "/requestURL";
+        String url = BASE_URL + "/api/power/electrical-model/query";
         log.info("配网电气模型查询-请求参数:{}", JSON.toJSONString(request));
 
         HttpHeaders headers = new HttpHeaders();
@@ -228,9 +228,9 @@ public class SimulationServiceImpl implements ISimulationService {
 
         HttpEntity<String> requestEntity = new HttpEntity<>(JSON.toJSONString(request), headers);
         try {
-            ElectricalModelResponse rs  = simpleRestTemplate.postForObject(url, requestEntity, ElectricalModelResponse.class);
-            log.info("配网电气模型查询-响应:{}", JSON.toJSONString(rs));
-            return rs;
+            String responseStr = simpleRestTemplate.postForObject(url, requestEntity, String.class);
+            log.info("配网电气模型查询-响应:{}", responseStr);
+            return responseStr;
         } catch (Exception e) {
             log.error("配网电气模型查询失败", e);
             ErrorLogParam errorLogParam = new ErrorLogParam()