lisonglin 3 недель назад
Родитель
Сommit
5a40006ceb
40 измененных файлов с 2303 добавлено и 0 удалено
  1. 50 0
      api/load-transfer-mq-api/src/main/java/com/hdkj/lt/modle/dto/indicator/IndicatorCalcMessage.java
  2. 86 0
      api/load-transfer-mq-api/src/main/java/com/hdkj/lt/modle/dto/indicator/IndicatorCompletedEvent.java
  3. 28 0
      api/load-transfer-si-api/src/main/java/com/hdkj/lt/feign/IAlgorithmApiClient.java
  4. 6 0
      common/common-core/src/main/java/com/hdkj/lt/base/constants/RedisKeyConstants.java
  5. 21 0
      common/common-core/src/main/java/com/hdkj/lt/base/constants/RocketMQConstants.java
  6. 42 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/controller/optimization/AlarmController.java
  7. 105 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/AlarmEvent.java
  8. 114 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/IndicatorResult.java
  9. 25 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/dto/AlarmIndicatorDTO.java
  10. 78 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/vo/AlarmFeederVO.java
  11. 32 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/mapper/AlarmEventMapper.java
  12. 12 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/mapper/IndicatorResultMapper.java
  13. 86 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/scheduler/IndicatorScheduler.java
  14. 22 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/AlarmEventService.java
  15. 19 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/IndicatorResultService.java
  16. 145 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/AlarmEventServiceImpl.java
  17. 25 0
      services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/IndicatorResultServiceImpl.java
  18. 123 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/entity/indicator/AlarmEvent.java
  19. 172 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/entity/indicator/IndicatorResult.java
  20. 128 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/listener/IndicatorCalcConsumer.java
  21. 139 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/listener/IndicatorEventConsumer.java
  22. 15 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/mapper/indicator/AlarmEventMapper.java
  23. 15 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/mapper/indicator/IndicatorResultMapper.java
  24. 28 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/AlarmEvaluateService.java
  25. 31 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/IndicatorResultService.java
  26. 249 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/impl/AlarmEvaluateServiceImpl.java
  27. 175 0
      services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/impl/IndicatorResultServiceImpl.java
  28. 33 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/controller/algorithm/AlgorithmController.java
  29. 25 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/Alarm.java
  30. 19 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/AlgorithmTrigger.java
  31. 15 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/CountyResult.java
  32. 17 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/FeederResult.java
  33. 44 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/IndicatorStatusResponse.java
  34. 15 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/LoadRate.java
  35. 15 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/Loss.java
  36. 18 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/ReconfigurationFeeders.java
  37. 14 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/VoltagePu.java
  38. 16 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/algorithm/AlgorithmService.java
  39. 27 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/algorithm/impl/AlgorithmServiceImpl.java
  40. 74 0
      services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/remote/impl/AlgorithmRemoteCallServiceImpl.java

+ 50 - 0
api/load-transfer-mq-api/src/main/java/com/hdkj/lt/modle/dto/indicator/IndicatorCalcMessage.java

@@ -0,0 +1,50 @@
+package com.hdkj.lt.modle.dto.indicator;
+
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.io.Serializable;
+import java.util.Map;
+
+/**
+ * 指标计算任务消息体(indicator-calc-topic)
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+@Builder
+public class IndicatorCalcMessage implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    /**
+     * 批次ID
+     */
+    private String batchId;
+
+    /**
+     * 区县ID
+     */
+    private Integer countyId;
+
+    /**
+     * 区县短标识
+     */
+    private String countySid;
+
+    /**
+     * 区县名称
+     */
+    private String countyName;
+
+    /**
+     * 扩展参数(预留,断面推送上线后填充断面数据)
+     */
+    @Builder.Default
+    private Map<String, Object> extParams = new java.util.HashMap<>();
+}

+ 86 - 0
api/load-transfer-mq-api/src/main/java/com/hdkj/lt/modle/dto/indicator/IndicatorCompletedEvent.java

@@ -0,0 +1,86 @@
+package com.hdkj.lt.modle.dto.indicator;
+
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.io.Serializable;
+import java.util.List;
+
+/**
+ * 指标完成事件(indicator-event-topic)
+ * <p>
+ * 事件带决策数据(告警标志+超限馈线),不带展示数据(负载率数值等在库里)
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+@Builder
+public class IndicatorCompletedEvent implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    /**
+     * 批次ID
+     */
+    private String batchId;
+
+    /**
+     * 区县ID
+     */
+    private Integer countyId;
+
+    /**
+     * 区县短标识
+     */
+    private String countySid;
+
+    /**
+     * 区县名称
+     */
+    private String countyName;
+
+    /**
+     * 指标结果ID(t_indicator_result.id)
+     */
+    private Long indicatorResultId;
+
+    /**
+     * 负载率超限 0/1
+     */
+    private Integer loadRateOverLimit;
+
+    /**
+     * 电压越上限 0/1
+     */
+    private Integer voltageOverUpperLimit;
+
+    /**
+     * 电压越下限 0/1
+     */
+    private Integer voltageUnderLowerLimit;
+
+    /**
+     * 超限馈线ID列表
+     */
+    private List<String> overLimitFeederIds;
+
+    /**
+     * 算法是否触发重构 0/1
+     */
+    private Integer reconfigurationTriggered;
+
+    /**
+     * 触发类型
+     */
+    private String triggerType;
+
+    /**
+     * 问题馈线ID列表
+     */
+    private List<String> problemFeederIds;
+}

+ 28 - 0
api/load-transfer-si-api/src/main/java/com/hdkj/lt/feign/IAlgorithmApiClient.java

@@ -0,0 +1,28 @@
+package com.hdkj.lt.feign;
+
+import com.hdkj.lt.base.constants.SystemGlobalConstant;
+import org.springframework.cloud.openfeign.FeignClient;
+import org.springframework.web.bind.annotation.PostMapping;
+
+/**
+ * 算法服务 Feign 客户端
+ * <p>
+ * 供 mq/bf 模块通过 Feign 调用 si 模块的算法指标计算接口
+ * <p>
+ * 返回 ApiResponse 完整 JSON 字符串,消费方用 fastjson 解析取 data 字段
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@FeignClient(path = "/algorithm", value = SystemGlobalConstant.SERVICE_NAME_LOAD_TRANSFER_SI,
+        contextId = "IAlgorithmApiClient")
+public interface IAlgorithmApiClient {
+
+    /**
+     * 查询指标计算结果
+     *
+     * @return ApiResponse 完整 JSON 字符串
+     */
+    @PostMapping("/indicators")
+    String queryIndicators();
+}

+ 6 - 0
common/common-core/src/main/java/com/hdkj/lt/base/constants/RedisKeyConstants.java

@@ -35,6 +35,12 @@ public class RedisKeyConstants {
 
     public static final String ASYNC_QUEUE_SCHEME_SCORE_TICKET = "schemeScoreTicket";
 
+    // 指标计算异步队列
+    public static final String ASYNC_QUEUE_INDICATOR_CALC = "indicatorCalcTicket";
+
+    // 指标完成事件异步队列
+    public static final String ASYNC_QUEUE_INDICATOR_EVENT = "indicatorEventTicket";
+
     // 异步处理队列 end
 
     // 组织机构

+ 21 - 0
common/common-core/src/main/java/com/hdkj/lt/base/constants/RocketMQConstants.java

@@ -43,5 +43,26 @@ public final class RocketMQConstants {
      * 处理配电线路动态映射表数据
      */
     public static final String TOPIC_ID_DIST_LINE_MAPPING_FILE_TICKET = "dist_line-mapping-file-topic";
+
+    /**
+     * 指标计算任务
+     */
+    public static final String TOPIC_ID_INDICATOR_CALC = "indicator-calc-topic";
+
+    /**
+     * 指标完成事件
+     */
+    public static final String TOPIC_ID_INDICATOR_EVENT = "indicator-event-topic";
+
+    /**
+     * 重构计算任务(预留)
+     */
+    public static final String TOPIC_ID_RECONFIGURATION_CALC = "reconfiguration-calc-topic";
+
+    /**
+     * 重构完成事件(预留)
+     */
+    public static final String TOPIC_ID_RECONFIGURATION_EVENT = "reconfiguration-event-topic";
+
     // topic-id end
 }

+ 42 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/controller/optimization/AlarmController.java

@@ -0,0 +1,42 @@
+package com.hdkj.lt.bf.controller.optimization;
+
+import com.hdkj.hussar.ApiResponse;
+import com.hdkj.lt.bf.entity.vo.AlarmFeederVO;
+import com.hdkj.lt.bf.service.AlarmEventService;
+import com.hdkj.lt.core.mvc.BaseController;
+import lombok.AllArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+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.time.LocalDate;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 告警 Controller
+ */
+@Slf4j
+@RestController
+@AllArgsConstructor
+@RequestMapping("/alarm")
+public class AlarmController extends BaseController {
+
+    private final AlarmEventService alarmEventService;
+
+    /**
+     * 查询馈线告警列表(联表 t_psr 获取馈线名称,返回中文 VO)
+     *
+     * @param params 请求参数,含 date(yyyy-MM-dd)
+     * @return 告警 VO 列表
+     */
+    @PostMapping("/feeders")
+    public ApiResponse<List<AlarmFeederVO>> feeders(@RequestBody Map<String, String> params) {
+        String date = params.get("date");
+        LocalDate alarmDate = LocalDate.parse(date);
+        List<AlarmFeederVO> list = alarmEventService.getAlarmFeederList(alarmDate);
+        return ApiResponse.success(list);
+    }
+}

+ 105 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/AlarmEvent.java

@@ -0,0 +1,105 @@
+package com.hdkj.lt.bf.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableField;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import com.fasterxml.jackson.annotation.JsonFormat;
+import io.swagger.annotations.ApiModel;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.io.Serializable;
+import java.math.BigDecimal;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+
+/**
+ * 告警事件实体
+ */
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+@Builder
+@TableName("t_alarm_event")
+@ApiModel(description = "告警事件")
+public class AlarmEvent implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    @ApiModelProperty(value = "主键ID")
+    @TableId(type = IdType.AUTO)
+    private Long id;
+
+    @ApiModelProperty(value = "批次ID")
+    private String batchId;
+
+    @ApiModelProperty(value = "区县ID")
+    private String countyId;
+
+    @ApiModelProperty(value = "区县SID")
+    private String countySid;
+
+    @ApiModelProperty(value = "区县名称")
+    private String countyName;
+
+    @ApiModelProperty(value = "馈线ID")
+    private String feederId;
+
+    @ApiModelProperty(value = "告警类型")
+    private String alarmTypes;
+
+    @ApiModelProperty(value = "最大负载率(%)")
+    private BigDecimal maxLoadRatePct;
+
+    @ApiModelProperty(value = "最大电压(pu)")
+    private BigDecimal maxVoltagePu;
+
+    @ApiModelProperty(value = "最小电压(pu)")
+    private BigDecimal minVoltagePu;
+
+    @ApiModelProperty(value = "指标结果ID")
+    private Long indicatorResultId;
+
+    @ApiModelProperty(value = "告警次数")
+    private Integer alarmCount;
+
+    @ApiModelProperty(value = "告警日期")
+    private LocalDate alarmDate;
+
+    @ApiModelProperty(value = "告警时间")
+    @JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    private LocalDateTime alarmTime;
+
+    @ApiModelProperty(value = "最后告警时间")
+    @JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    private LocalDateTime lastAlarmTime;
+
+    @ApiModelProperty(value = "状态")
+    private Integer status;
+
+    @ApiModelProperty(value = "创建时间")
+    @JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    private LocalDateTime createTime;
+
+    @ApiModelProperty(value = "更新时间")
+    @JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    private LocalDateTime updateTime;
+
+    /**
+     * 馈线名称(联表查询,非数据库字段)
+     */
+    @TableField(exist = false)
+    @ApiModelProperty(value = "馈线名称")
+    private String feederName;
+
+    /**
+     * 健康度(联表查询,非数据库字段)
+     */
+    @TableField(exist = false)
+    @ApiModelProperty(value = "健康度")
+    private String healthScore;
+}

+ 114 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/IndicatorResult.java

@@ -0,0 +1,114 @@
+package com.hdkj.lt.bf.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import io.swagger.annotations.ApiModel;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.io.Serializable;
+import java.math.BigDecimal;
+import java.time.LocalDateTime;
+
+/**
+ * 指标计算结果实体
+ */
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+@Builder
+@TableName("t_indicator_result")
+@ApiModel(description = "指标计算结果")
+public class IndicatorResult implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    @ApiModelProperty(value = "主键ID")
+    @TableId(type = IdType.AUTO)
+    private Long id;
+
+    @ApiModelProperty(value = "批次ID")
+    private String batchId;
+
+    @ApiModelProperty(value = "区县ID")
+    private Integer countyId;
+
+    @ApiModelProperty(value = "区县SID")
+    private String countySid;
+
+    @ApiModelProperty(value = "区县名称")
+    private String countyName;
+
+    @ApiModelProperty(value = "状态")
+    private String status;
+
+    @ApiModelProperty(value = "计算结果")
+    private String computeResult;
+
+    @ApiModelProperty(value = "生成时间")
+    private String generatedAt;
+
+    @ApiModelProperty(value = "最大负载率(%)")
+    private BigDecimal maxLoadRatePct;
+
+    @ApiModelProperty(value = "最小负载率(%)")
+    private BigDecimal minLoadRatePct;
+
+    @ApiModelProperty(value = "平均负载率(%)")
+    private BigDecimal avgLoadRatePct;
+
+    @ApiModelProperty(value = "告警阈值(%)")
+    private BigDecimal alarmThresholdPct;
+
+    @ApiModelProperty(value = "最大电压(pu)")
+    private BigDecimal maxVoltagePu;
+
+    @ApiModelProperty(value = "最小电压(pu)")
+    private BigDecimal minVoltagePu;
+
+    @ApiModelProperty(value = "平均电压(pu)")
+    private BigDecimal avgVoltagePu;
+
+    @ApiModelProperty(value = "总损耗率(%)")
+    private BigDecimal totalLossRatePct;
+
+    @ApiModelProperty(value = "变压器损耗率(%)")
+    private BigDecimal transformerLossRatePct;
+
+    @ApiModelProperty(value = "线路损耗率(%)")
+    private BigDecimal lineLossRatePct;
+
+    @ApiModelProperty(value = "负载率越限")
+    private Integer loadRateOverLimit;
+
+    @ApiModelProperty(value = "电压越上限")
+    private Integer voltageOverUpperLimit;
+
+    @ApiModelProperty(value = "电压越下限")
+    private Integer voltageUnderLowerLimit;
+
+    @ApiModelProperty(value = "越限数量")
+    private Integer overLimitCount;
+
+    @ApiModelProperty(value = "越限馈线ID")
+    private String overLimitFeederIds;
+
+    @ApiModelProperty(value = "是否触发重构")
+    private Integer reconfigurationTriggered;
+
+    @ApiModelProperty(value = "触发类型")
+    private String triggerType;
+
+    @ApiModelProperty(value = "问题馈线ID")
+    private String problemFeederIds;
+
+    @ApiModelProperty(value = "原始响应")
+    private String rawResponse;
+
+    @ApiModelProperty(value = "创建时间")
+    private LocalDateTime createTime;
+}

+ 25 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/dto/AlarmIndicatorDTO.java

@@ -0,0 +1,25 @@
+package com.hdkj.lt.bf.entity.dto;
+
+import com.fasterxml.jackson.databind.annotation.JsonSerialize;
+import com.hdkj.lt.base.serializer.BigDecimalSerializer;
+import lombok.AllArgsConstructor;
+import lombok.Data;
+
+import java.math.BigDecimal;
+
+/**
+ * @author lisonglin
+ * @date: 2026/7/18
+ * @Description
+ */
+@Data
+@AllArgsConstructor
+public class AlarmIndicatorDTO {
+
+    private String label;
+
+    @JsonSerialize(using = BigDecimalSerializer.StripTrailingZerosSerializer.class)
+    private BigDecimal value;
+
+
+}

+ 78 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/entity/vo/AlarmFeederVO.java

@@ -0,0 +1,78 @@
+package com.hdkj.lt.bf.entity.vo;
+
+import com.fasterxml.jackson.annotation.JsonFormat;
+import com.fasterxml.jackson.databind.annotation.JsonSerialize;
+import com.hdkj.lt.base.serializer.BigDecimalSerializer;
+import com.hdkj.lt.bf.entity.dto.AlarmIndicatorDTO;
+import io.swagger.annotations.ApiModel;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.io.Serializable;
+import java.math.BigDecimal;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.util.List;
+
+/**
+ * 馈线告警列表 VO(前端展示用,中文化)
+ */
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+@Builder
+@ApiModel(description = "馈线告警列表")
+public class AlarmFeederVO implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    @ApiModelProperty(value = "告警事件ID")
+    private Long id;
+
+    @ApiModelProperty(value = "所属区县")
+    private String countyName;
+
+    @ApiModelProperty(value = "异常线路")
+    private String feederId;
+
+    @ApiModelProperty(value = "异常线路")
+    private String feederName;
+
+    @ApiModelProperty(value = "异常指标标签(中文,红色标签展示)")
+    private List<String> alarmLabels;
+
+    @ApiModelProperty(value = "异常指标(中文标签+数值拼接)")
+    private List<AlarmIndicatorDTO> alarmIndicators;
+
+    @ApiModelProperty(value = "最大负载率(%)")
+    @JsonSerialize(using = BigDecimalSerializer.StripTrailingZerosSerializer.class)
+    private BigDecimal maxLoadRatePct;
+
+    @ApiModelProperty(value = "告警次数")
+    private Integer alarmCount;
+
+    @ApiModelProperty(value = "告警日期")
+    private LocalDate alarmDate;
+
+    @ApiModelProperty(value = "首次告警时间")
+    @JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    private LocalDateTime alarmTime;
+
+    @ApiModelProperty(value = "最近告警时间")
+    @JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    private LocalDateTime lastAlarmTime;
+
+    @ApiModelProperty(value = "状态(中文)")
+    private String statusName;
+
+    @ApiModelProperty(value = "关联指标结果ID")
+    private Long indicatorResultId;
+
+    @ApiModelProperty(value = "健康度")
+    private String healthScore;
+
+
+}

+ 32 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/mapper/AlarmEventMapper.java

@@ -0,0 +1,32 @@
+package com.hdkj.lt.bf.mapper;
+
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import com.hdkj.lt.bf.entity.AlarmEvent;
+import org.apache.ibatis.annotations.Mapper;
+import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.annotations.Select;
+
+import java.time.LocalDate;
+import java.util.List;
+
+/**
+ * 告警事件 Mapper
+ */
+@Mapper
+public interface AlarmEventMapper extends BaseMapper<AlarmEvent> {
+
+    /**
+     * 按日期查询告警列表,联表 t_psr 获取馈线名称
+     *
+     * @param alarmDate 告警日期
+     * @return 告警列表(含馈线名称)
+     */
+    @Select("SELECT a.*, p.name AS feeder_name, rm.before_health_score AS healthScore " +
+            "FROM t_alarm_event a " +
+            "LEFT JOIN dwd_shb_ds_feeder_base p ON a.feeder_id = p.psr_id " +
+            "LEFT JOIN fhzg_recon_result rr ON rr.event_id = a.id " +
+            "LEFT JOIN fhzg_recon_metrics rm ON rm.result_id = rr.id " +
+            "WHERE a.alarm_date = #{alarmDate} " +
+            "ORDER BY a.last_alarm_time DESC")
+    List<AlarmEvent> selectAlarmListWithFeederName(@Param("alarmDate") LocalDate alarmDate);
+}

+ 12 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/mapper/IndicatorResultMapper.java

@@ -0,0 +1,12 @@
+package com.hdkj.lt.bf.mapper;
+
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import com.hdkj.lt.bf.entity.IndicatorResult;
+import org.apache.ibatis.annotations.Mapper;
+
+/**
+ * 指标计算结果 Mapper
+ */
+@Mapper
+public interface IndicatorResultMapper extends BaseMapper<IndicatorResult> {
+}

+ 86 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/scheduler/IndicatorScheduler.java

@@ -0,0 +1,86 @@
+package com.hdkj.lt.bf.scheduler;
+
+import com.hdkj.lt.base.constants.RedisKeyConstants;
+import com.hdkj.lt.base.constants.RocketMQConstants;
+import com.hdkj.lt.core.mq.client.RocketMQTemplate;
+import com.hdkj.lt.core.mq.model.RocketMQMessageEntity;
+import com.hdkj.lt.modle.dto.indicator.IndicatorCalcMessage;
+import lombok.Getter;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+import java.time.LocalDateTime;
+import java.time.format.DateTimeFormatter;
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * 指标计算定时调度器
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class IndicatorScheduler {
+
+    private final RocketMQTemplate rocketMQTemplate;
+
+    /**
+     * 每15分钟触发一次指标计算
+     */
+    //@Scheduled(cron = "0 */15 * * * ?")
+    public void executeIndicatorCalc() {
+        String batchId = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyyMMddHHmmssSSS"));
+        log.info("指标计算定时任务启动, batchId: {}", batchId);
+
+        List<CountyInfo> countyList = getCountyList();
+        for (CountyInfo county : countyList) {
+            try {
+                IndicatorCalcMessage message = new IndicatorCalcMessage();
+                message.setBatchId(batchId);
+                message.setCountyId(county.getId());
+                message.setCountySid(county.getSid());
+                message.setCountyName(county.getName());
+
+                RocketMQMessageEntity messageEntity = RocketMQMessageEntity.builder()
+                        .topicId(RocketMQConstants.TOPIC_ID_INDICATOR_CALC)
+                        .cacheKey(RedisKeyConstants.ASYNC_QUEUE_INDICATOR_CALC)
+                        .tag(county.getSid())
+                        .body(message)
+                        .build();
+                rocketMQTemplate.sendAsyncMessage(messageEntity);
+                log.info("发送指标计算消息成功, countySid: {}, batchId: {}", county.getSid(), batchId);
+            } catch (Exception e) {
+                log.error("发送指标计算消息失败, countySid: {}, batchId: {}", county.getSid(), batchId, e);
+            }
+        }
+        log.info("指标计算定时任务结束, batchId: {}", batchId);
+    }
+
+    /**
+     * 获取区县列表 (TODO 待实现)
+     *
+     * @return 区县列表
+     */
+    private List<CountyInfo> getCountyList() {
+        // TODO 待实现:从数据库或其他数据源获取区县列表
+        return new ArrayList<>();
+    }
+
+    /**
+     * 区县信息
+     */
+    @Getter
+    public static class CountyInfo {
+        private Integer id;
+        private String sid;
+        private String name;
+
+        public CountyInfo(Integer id, String sid, String name) {
+            this.id = id;
+            this.sid = sid;
+            this.name = name;
+        }
+    }
+}

+ 22 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/AlarmEventService.java

@@ -0,0 +1,22 @@
+package com.hdkj.lt.bf.service;
+
+import com.baomidou.mybatisplus.extension.service.IService;
+import com.hdkj.lt.bf.entity.AlarmEvent;
+import com.hdkj.lt.bf.entity.vo.AlarmFeederVO;
+
+import java.time.LocalDate;
+import java.util.List;
+
+/**
+ * 告警事件 Service
+ */
+public interface AlarmEventService extends IService<AlarmEvent> {
+
+    /**
+     * 按日期查询告警列表(联表 t_psr 获取馈线名称,转换为中文 VO)
+     *
+     * @param alarmDate 告警日期
+     * @return 告警 VO 列表
+     */
+    List<AlarmFeederVO> getAlarmFeederList(LocalDate alarmDate);
+}

+ 19 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/IndicatorResultService.java

@@ -0,0 +1,19 @@
+package com.hdkj.lt.bf.service;
+
+import com.baomidou.mybatisplus.extension.service.IService;
+import com.hdkj.lt.bf.entity.IndicatorResult;
+
+/**
+ * 指标计算结果 Service
+ */
+public interface IndicatorResultService extends IService<IndicatorResult> {
+
+    /**
+     * 根据批次ID和区县ID查询指标结果
+     *
+     * @param batchId  批次ID
+     * @param countyId 区县ID
+     * @return 指标结果
+     */
+    IndicatorResult getByBatchAndCounty(String batchId, Integer countyId);
+}

+ 145 - 0
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/AlarmEventServiceImpl.java

@@ -0,0 +1,145 @@
+package com.hdkj.lt.bf.service.impl;
+
+import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
+import com.hdkj.lt.bf.entity.AlarmEvent;
+import com.hdkj.lt.bf.entity.dto.AlarmIndicatorDTO;
+import com.hdkj.lt.bf.entity.vo.AlarmFeederVO;
+import com.hdkj.lt.bf.mapper.AlarmEventMapper;
+import com.hdkj.lt.bf.service.AlarmEventService;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Service;
+
+import java.math.BigDecimal;
+import java.time.LocalDate;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * 告警事件 Service 实现
+ */
+@Slf4j
+@Service
+public class AlarmEventServiceImpl extends ServiceImpl<AlarmEventMapper, AlarmEvent> implements AlarmEventService {
+
+    /**
+     * status → 中文
+     */
+    private static final Map<Integer, String> STATUS_LABELS = new LinkedHashMap<Integer, String>();
+    static {
+        STATUS_LABELS.put(0, "待处理");
+        STATUS_LABELS.put(1, "重构中");
+        STATUS_LABELS.put(2, "已完成");
+        STATUS_LABELS.put(3, "已忽略");
+    }
+
+    @Override
+    public List<AlarmFeederVO> getAlarmFeederList(LocalDate alarmDate) {
+        List<AlarmEvent> list = baseMapper.selectAlarmListWithFeederName(alarmDate);
+        List<AlarmFeederVO> result = new ArrayList<>(list.size());
+        for (AlarmEvent event : list) {
+            AlarmFeederVO vo = AlarmFeederVO.builder()
+                    .id(event.getId())
+                    .countyName(event.getCountyName())
+                    .feederId(event.getFeederId())
+                    .feederName(event.getFeederName())
+                    .alarmLabels(convertAlarmLabels(event.getAlarmTypes()))
+                    .alarmIndicators(buildAlarmIndicators(event))
+                    .alarmCount(event.getAlarmCount())
+                    .alarmDate(event.getAlarmDate())
+                    .alarmTime(event.getAlarmTime())
+                    .lastAlarmTime(event.getLastAlarmTime())
+                    .statusName(convertStatus(event.getStatus()))
+                    .indicatorResultId(event.getIndicatorResultId())
+                    .maxLoadRatePct(event.getMaxLoadRatePct())
+                    .healthScore(event.getHealthScore())
+                    .build();
+            result.add(vo);
+        }
+        return result;
+    }
+
+    /**
+     * alarm_types 英文枚举 → 中文标签列表(去重),用于红色标签展示
+     */
+    private List<String> convertAlarmLabels(String alarmTypes) {
+        List<String> labels = new ArrayList<String>();
+        if (StringUtils.isBlank(alarmTypes)) {
+            return labels;
+        }
+        Set<String> seen = new LinkedHashSet<String>();
+        for (String type : alarmTypes.split(",")) {
+            String trimmed = type.trim();
+            String label = null;
+            if ("load_rate".equals(trimmed)) {
+                label = "线路重过载";
+            } else if ("voltage_over".equals(trimmed) || "voltage_under".equals(trimmed)) {
+                label = "电压越限";
+            }
+            if (label != null && seen.add(label)) {
+                labels.add(label);
+            }
+        }
+        return labels;
+    }
+
+    /**
+     * 构建"异常指标"展示字符串列表
+     * <p>
+     * 根据 alarm_types 拼接中文标签 + 数值:
+     * <ul>
+     *   <li>load_rate    → "线路重过载 79.70%"</li>
+     *   <li>voltage_over  → "电压越限 1.05pu"</li>
+     *   <li>voltage_under → "电压越限 0.92pu"</li>
+     * </ul>
+     * voltage_over 和 voltage_under 都映射"电压越限",但数值不同:
+     * - voltage_over 取 maxVoltagePu(越上限,取最大值)
+     * - voltage_under 取 minVoltagePu(越下限,取最小值)
+     */
+    private List<AlarmIndicatorDTO> buildAlarmIndicators(AlarmEvent event) {
+        List<AlarmIndicatorDTO> result = new ArrayList<AlarmIndicatorDTO>();
+        if (StringUtils.isBlank(event.getAlarmTypes())) {
+            return result;
+        }
+
+        Set<String> seen = new LinkedHashSet<>();
+        for (String type : event.getAlarmTypes().split(",")) {
+            String trimmed = type.trim();
+            if ("load_rate".equals(trimmed)) {
+                String label = "线路重过载";
+                if (seen.add(label)) {
+                    BigDecimal val = event.getMaxLoadRatePct();
+                    result.add(new AlarmIndicatorDTO(label, val));
+                }
+            } else if ("voltage_over".equals(trimmed)) {
+                String label = "电压越限";
+                if (seen.add(label)) {
+                    BigDecimal val = event.getMaxVoltagePu();
+                    result.add(new AlarmIndicatorDTO(label, val));
+                }
+            } else if ("voltage_under".equals(trimmed)) {
+                String label = "电压越限";
+                if (seen.add(label)) {
+                    BigDecimal val = event.getMinVoltagePu();
+                    result.add(new AlarmIndicatorDTO(label, val));
+                }
+            }
+        }
+        return result;
+    }
+
+    /**
+     * status 数字 → 中文
+     */
+    private String convertStatus(Integer status) {
+        if (status == null) {
+            return null;
+        }
+        return STATUS_LABELS.get(status);
+    }
+}

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

@@ -0,0 +1,25 @@
+package com.hdkj.lt.bf.service.impl;
+
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
+import com.hdkj.lt.bf.entity.IndicatorResult;
+import com.hdkj.lt.bf.mapper.IndicatorResultMapper;
+import com.hdkj.lt.bf.service.IndicatorResultService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+
+/**
+ * 指标计算结果 Service 实现
+ */
+@Slf4j
+@Service
+public class IndicatorResultServiceImpl extends ServiceImpl<IndicatorResultMapper, IndicatorResult> implements IndicatorResultService {
+
+    @Override
+    public IndicatorResult getByBatchAndCounty(String batchId, Integer countyId) {
+        LambdaQueryWrapper<IndicatorResult> queryWrapper = new LambdaQueryWrapper<IndicatorResult>()
+                .eq(IndicatorResult::getBatchId, batchId)
+                .eq(IndicatorResult::getCountyId, countyId);
+        return baseMapper.selectOne(queryWrapper);
+    }
+}

+ 123 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/entity/indicator/AlarmEvent.java

@@ -0,0 +1,123 @@
+package com.hdkj.lt.mq.entity.indicator;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableField;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.io.Serializable;
+import java.math.BigDecimal;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+
+/**
+ * 告警事件(t_alarm_event)
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@TableName("t_alarm_event")
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+@Builder
+public class AlarmEvent implements Serializable {
+
+    @TableField(exist = false)
+    private static final long serialVersionUID = 1L;
+
+    /**
+     * 主键
+     */
+    @TableId(type = IdType.AUTO)
+    private Long id;
+
+    /**
+     * 批次ID
+     */
+    private String batchId;
+
+    /**
+     * 区县ID
+     */
+    private Integer countyId;
+
+    /**
+     * 区县短标识
+     */
+    private String countySid;
+
+    /**
+     * 区县名称
+     */
+    private String countyName;
+
+    /**
+     * 馈线ID
+     */
+    private String feederId;
+
+    /**
+     * 告警类型(逗号分隔)
+     */
+    private String alarmTypes;
+
+    /**
+     * 最大负载率(百分比)
+     */
+    private BigDecimal maxLoadRatePct;
+
+    /**
+     * 最大电压(pu)
+     */
+    private BigDecimal maxVoltagePu;
+
+    /**
+     * 最小电压(pu)
+     */
+    private BigDecimal minVoltagePu;
+
+    /**
+     * 关联指标结果ID
+     */
+    private Long indicatorResultId;
+
+    /**
+     * 告警次数(累计)
+     */
+    private Integer alarmCount;
+
+    /**
+     * 告警日期
+     */
+    private LocalDate alarmDate;
+
+    /**
+     * 首次告警时间
+     */
+    private LocalDateTime alarmTime;
+
+    /**
+     * 最近告警时间
+     */
+    private LocalDateTime lastAlarmTime;
+
+    /**
+     * 状态:0-待处理,1-重构中,2-已完成,3-已忽略
+     */
+    private Integer status;
+
+    /**
+     * 创建时间
+     */
+    private LocalDateTime createTime;
+
+    /**
+     * 更新时间
+     */
+    private LocalDateTime updateTime;
+}

+ 172 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/entity/indicator/IndicatorResult.java

@@ -0,0 +1,172 @@
+package com.hdkj.lt.mq.entity.indicator;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableField;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.io.Serializable;
+import java.math.BigDecimal;
+import java.time.LocalDateTime;
+
+/**
+ * 指标计算结果(t_indicator_result)
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@TableName("t_indicator_result")
+@Data
+@NoArgsConstructor
+@AllArgsConstructor
+@Builder
+public class IndicatorResult implements Serializable {
+
+    @TableField(exist = false)
+    private static final long serialVersionUID = 1L;
+
+    /**
+     * 主键
+     */
+    @TableId(type = IdType.AUTO)
+    private Long id;
+
+    /**
+     * 批次ID
+     */
+    private String batchId;
+
+    /**
+     * 区县ID
+     */
+    private Integer countyId;
+
+    /**
+     * 区县短标识
+     */
+    private String countySid;
+
+    /**
+     * 区县名称
+     */
+    private String countyName;
+
+    /**
+     * 计算状态
+     */
+    private String status;
+
+    /**
+     * 计算结果
+     */
+    private String computeResult;
+
+    /**
+     * 生成时间(算法返回)
+     */
+    private String generatedAt;
+
+    /**
+     * 最大负载率(百分比)
+     */
+    private BigDecimal maxLoadRatePct;
+
+    /**
+     * 最小负载率(百分比)
+     */
+    private BigDecimal minLoadRatePct;
+
+    /**
+     * 平均负载率(百分比)
+     */
+    private BigDecimal avgLoadRatePct;
+
+    /**
+     * 告警阈值(百分比)
+     */
+    private BigDecimal alarmThresholdPct;
+
+    /**
+     * 最大电压(pu)
+     */
+    private BigDecimal maxVoltagePu;
+
+    /**
+     * 最小电压(pu)
+     */
+    private BigDecimal minVoltagePu;
+
+    /**
+     * 平均电压(pu)
+     */
+    private BigDecimal avgVoltagePu;
+
+    /**
+     * 总损耗率(百分比)
+     */
+    private BigDecimal totalLossRatePct;
+
+    /**
+     * 变压器损耗率(百分比)
+     */
+    private BigDecimal transformerLossRatePct;
+
+    /**
+     * 线路损耗率(百分比)
+     */
+    private BigDecimal lineLossRatePct;
+
+    /**
+     * 负载率超限 0/1
+     */
+    private Integer loadRateOverLimit;
+
+    /**
+     * 电压越上限 0/1
+     */
+    private Integer voltageOverUpperLimit;
+
+    /**
+     * 电压越下限 0/1
+     */
+    private Integer voltageUnderLowerLimit;
+
+    /**
+     * 超限数量
+     */
+    private Integer overLimitCount;
+
+    /**
+     * 超限馈线ID列表(逗号分隔)
+     */
+    private String overLimitFeederIds;
+
+    /**
+     * 算法是否触发重构 0/1
+     */
+    private Integer reconfigurationTriggered;
+
+    /**
+     * 触发类型
+     */
+    private String triggerType;
+
+    /**
+     * 问题馈线ID列表(逗号分隔)
+     */
+    private String problemFeederIds;
+
+    /**
+     * 原始响应JSON
+     */
+    private String rawResponse;
+
+    /**
+     * 创建时间
+     */
+    private LocalDateTime createTime;
+}

+ 128 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/listener/IndicatorCalcConsumer.java

@@ -0,0 +1,128 @@
+package com.hdkj.lt.mq.listener;
+
+import com.alibaba.fastjson.JSON;
+import com.aliyun.openservices.ons.api.Action;
+import com.aliyun.openservices.ons.api.ConsumeContext;
+import com.aliyun.openservices.ons.api.Message;
+import com.aliyun.openservices.ons.api.MessageListener;
+import com.hdkj.lt.base.constants.RedisKeyConstants;
+import com.hdkj.lt.base.constants.RocketMQConstants;
+import com.hdkj.lt.core.mq.client.RocketMQTemplate;
+import com.hdkj.lt.core.mq.model.RocketMQMessageEntity;
+import com.hdkj.lt.feign.IAlgorithmApiClient;
+import com.hdkj.lt.modle.dto.indicator.IndicatorCalcMessage;
+import com.hdkj.lt.modle.dto.indicator.IndicatorCompletedEvent;
+import com.hdkj.lt.mq.entity.indicator.IndicatorResult;
+import com.hdkj.lt.mq.service.indicator.IndicatorResultService;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Component;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+/**
+ * 指标计算消费者(indicator-calc-topic)
+ * <p>
+ * 消费指标计算任务消息,通过 Feign 调 si 模块算法接口获取指标JSON,
+ * 入库后发送指标完成事件到 indicator-event-topic。
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class IndicatorCalcConsumer implements MessageListener {
+
+    private final IAlgorithmApiClient algorithmApiClient;
+    private final IndicatorResultService indicatorResultService;
+    private final RocketMQTemplate rocketMQTemplate;
+
+    @Override
+    public Action consume(Message message, ConsumeContext consumeContext) {
+        try {
+            log.info("===== 指标计算任务-开始消费 =====");
+            String body = new String(message.getBody(), StandardCharsets.UTF_8);
+            IndicatorCalcMessage calcMessage = JSON.parseObject(body, IndicatorCalcMessage.class);
+            if (calcMessage == null) {
+                log.warn("指标计算任务消息解析失败,body={}", body);
+                return Action.ReconsumeLater;
+            }
+            String batchId = calcMessage.getBatchId();
+            Integer countyId = calcMessage.getCountyId();
+
+            // 调算法接口获取指标JSON(返回 ApiResponse 外壳,data 字段为算法原始JSON)
+            String apiResponseJson = algorithmApiClient.queryIndicators();
+            if (StringUtils.isBlank(apiResponseJson)) {
+                log.warn("算法接口返回空 batchId={}, countyId={}", batchId, countyId);
+                return Action.ReconsumeLater;
+            }
+
+            // 从 ApiResponse 外壳提取 data(算法原始JSON)
+            com.alibaba.fastjson.JSONObject apiResp = JSON.parseObject(apiResponseJson);
+            String indicatorJson = apiResp.getString("data");
+            if (StringUtils.isBlank(indicatorJson)) {
+                log.warn("算法接口返回 data 为空, apiResponse={}", apiResponseJson);
+                return Action.ReconsumeLater;
+            }
+
+            // 入库
+            IndicatorResult result = indicatorResultService.saveResult(indicatorJson, batchId, countyId);
+
+            // 构建指标完成事件
+            IndicatorCompletedEvent event = IndicatorCompletedEvent.builder()
+                    .batchId(result.getBatchId())
+                    .countyId(result.getCountyId())
+                    .countySid(result.getCountySid())
+                    .countyName(result.getCountyName())
+                    .indicatorResultId(result.getId())
+                    .loadRateOverLimit(result.getLoadRateOverLimit())
+                    .voltageOverUpperLimit(result.getVoltageOverUpperLimit())
+                    .voltageUnderLowerLimit(result.getVoltageUnderLowerLimit())
+                    .overLimitFeederIds(splitToList(result.getOverLimitFeederIds()))
+                    .reconfigurationTriggered(result.getReconfigurationTriggered())
+                    .triggerType(result.getTriggerType())
+                    .problemFeederIds(splitToList(result.getProblemFeederIds()))
+                    .build();
+
+            // 发送到指标完成事件 topic
+            RocketMQMessageEntity messageEntity = RocketMQMessageEntity.builder()
+                    .topicId(RocketMQConstants.TOPIC_ID_INDICATOR_EVENT)
+                    .cacheKey(RedisKeyConstants.ASYNC_QUEUE_INDICATOR_EVENT)
+                    .tag(batchId)
+                    .body(event)
+                    .build();
+            rocketMQTemplate.sendAsyncMessage(messageEntity);
+
+            log.info("指标计算任务消费完成 batchId={}, countyId={}, indicatorResultId={}",
+                    batchId, countyId, result.getId());
+            return Action.CommitMessage;
+        } catch (Exception e) {
+            log.error("指标计算任务消费异常 message={}", message.getTopic(), e);
+            return Action.ReconsumeLater;
+        }
+    }
+
+    /**
+     * 逗号分隔字符串转 List
+     */
+    private List<String> splitToList(String commaStr) {
+        if (StringUtils.isBlank(commaStr)) {
+            return Collections.emptyList();
+        }
+        String[] arr = commaStr.split(",");
+        List<String> result = new ArrayList<String>(arr.length);
+        for (String s : arr) {
+            String trimmed = s.trim();
+            if (StringUtils.isNotBlank(trimmed)) {
+                result.add(trimmed);
+            }
+        }
+        return result;
+    }
+}

+ 139 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/listener/IndicatorEventConsumer.java

@@ -0,0 +1,139 @@
+package com.hdkj.lt.mq.listener;
+
+import com.alibaba.fastjson.JSON;
+import com.aliyun.openservices.ons.api.Action;
+import com.aliyun.openservices.ons.api.ConsumeContext;
+import com.aliyun.openservices.ons.api.Message;
+import com.aliyun.openservices.ons.api.MessageListener;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.hdkj.lt.modle.dto.indicator.IndicatorCompletedEvent;
+import com.hdkj.lt.mq.entity.indicator.AlarmEvent;
+import com.hdkj.lt.mq.entity.indicator.IndicatorResult;
+import com.hdkj.lt.mq.mapper.indicator.AlarmEventMapper;
+import com.hdkj.lt.mq.mapper.indicator.IndicatorResultMapper;
+import com.hdkj.lt.mq.service.indicator.AlarmEvaluateService;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Component;
+
+import java.nio.charset.StandardCharsets;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.util.Arrays;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Set;
+
+/**
+ * 指标完成事件消费者(indicator-event-topic)
+ * <p>
+ * 消费指标完成事件,触发告警评估,对告警事件做抑制写入(同馈线同天活跃告警合并)。
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class IndicatorEventConsumer implements MessageListener {
+
+    private final AlarmEvaluateService alarmEvaluateService;
+    private final AlarmEventMapper alarmEventMapper;
+    private final IndicatorResultMapper indicatorResultMapper;
+
+    @Override
+    public Action consume(Message message, ConsumeContext consumeContext) {
+        try {
+            log.info("===== 指标完成事件-开始消费 =====");
+            String body = new String(message.getBody(), StandardCharsets.UTF_8);
+            IndicatorCompletedEvent event = JSON.parseObject(body, IndicatorCompletedEvent.class);
+            if (event == null) {
+                log.warn("指标完成事件消息解析失败,body={}", body);
+                return Action.ReconsumeLater;
+            }
+            Long indicatorResultId = event.getIndicatorResultId();
+            String batchId = event.getBatchId();
+            Integer countyId = event.getCountyId();
+
+            // 查指标结果拿 rawResponse(算法原始JSON,含 over_limit_details + feeder_results)
+            IndicatorResult indicatorResult = indicatorResultMapper.selectById(indicatorResultId);
+            String rawResponseJson = indicatorResult != null ? indicatorResult.getRawResponse() : null;
+
+            // 告警评估
+            List<AlarmEvent> alarmEvents = alarmEvaluateService.evaluateAlarms(indicatorResultId, batchId, countyId, rawResponseJson);
+            if (alarmEvents == null || alarmEvents.isEmpty()) {
+                log.info("指标完成事件无告警 indicatorResultId={}, batchId={}, countyId={}", indicatorResultId, batchId, countyId);
+                return Action.CommitMessage;
+            }
+
+            // 遍历写入告警事件(抑制)
+            for (AlarmEvent alarmEvent : alarmEvents) {
+                upsertAlarmEvent(alarmEvent);
+            }
+
+            log.info("指标完成事件消费完成 indicatorResultId={}, batchId={}, countyId={}, alarmCount={}",
+                    indicatorResultId, batchId, countyId, alarmEvents.size());
+            return Action.CommitMessage;
+        } catch (Exception e) {
+            log.error("指标完成事件消费异常 message={}", message.getTopic(), e);
+            return Action.ReconsumeLater;
+        }
+    }
+
+    /**
+     * 告警抑制写入:同馈线同天活跃告警合并,否则新增
+     */
+    private void upsertAlarmEvent(AlarmEvent alarmEvent) {
+        LocalDate alarmDate = alarmEvent.getAlarmDate() != null ? alarmEvent.getAlarmDate() : LocalDate.now();
+        LambdaQueryWrapper<AlarmEvent> wrapper = new LambdaQueryWrapper<AlarmEvent>()
+                .eq(AlarmEvent::getCountyId, alarmEvent.getCountyId())
+                .eq(AlarmEvent::getFeederId, alarmEvent.getFeederId())
+                .eq(AlarmEvent::getAlarmDate, alarmDate)
+                .in(AlarmEvent::getStatus, 0, 1)
+                .last("LIMIT 1");
+        AlarmEvent existing = alarmEventMapper.selectOne(wrapper);
+
+        if (existing != null) {
+            // 合并更新
+            existing.setBatchId(alarmEvent.getBatchId());
+            existing.setAlarmTypes(mergeAlarmTypes(existing.getAlarmTypes(), alarmEvent.getAlarmTypes()));
+            existing.setMaxLoadRatePct(alarmEvent.getMaxLoadRatePct());
+            existing.setMaxVoltagePu(alarmEvent.getMaxVoltagePu());
+            existing.setMinVoltagePu(alarmEvent.getMinVoltagePu());
+            existing.setIndicatorResultId(alarmEvent.getIndicatorResultId());
+            existing.setAlarmCount(existing.getAlarmCount() != null ? existing.getAlarmCount() + 1 : 1);
+            existing.setLastAlarmTime(LocalDateTime.now());
+            existing.setUpdateTime(LocalDateTime.now());
+            alarmEventMapper.updateById(existing);
+            log.info("告警抑制-更新 alarmEventId={}, feederId={}, alarmCount={}",
+                    existing.getId(), existing.getFeederId(), existing.getAlarmCount());
+        } else {
+            // 新增
+            alarmEvent.setAlarmCount(1);
+            alarmEvent.setAlarmDate(alarmDate);
+            LocalDateTime now = LocalDateTime.now();
+            alarmEvent.setAlarmTime(now);
+            alarmEvent.setLastAlarmTime(now);
+            alarmEvent.setCreateTime(now);
+            alarmEvent.setUpdateTime(now);
+            alarmEvent.setStatus(0);
+            alarmEventMapper.insert(alarmEvent);
+            log.info("告警抑制-新增 feederId={}, countyId={}, alarmTypes={}",
+                    alarmEvent.getFeederId(), alarmEvent.getCountyId(), alarmEvent.getAlarmTypes());
+        }
+    }
+
+    /**
+     * 合并告警类型字符串,用 LinkedHashSet 去重后逗号拼接
+     */
+    private String mergeAlarmTypes(String existing, String newTypes) {
+        Set<String> set = new LinkedHashSet<String>();
+        if (existing != null && !existing.isEmpty()) {
+            set.addAll(Arrays.asList(existing.split(",")));
+        }
+        if (newTypes != null && !newTypes.isEmpty()) {
+            set.addAll(Arrays.asList(newTypes.split(",")));
+        }
+        return set.isEmpty() ? null : String.join(",", set);
+    }
+}

+ 15 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/mapper/indicator/AlarmEventMapper.java

@@ -0,0 +1,15 @@
+package com.hdkj.lt.mq.mapper.indicator;
+
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import com.hdkj.lt.mq.entity.indicator.AlarmEvent;
+import org.apache.ibatis.annotations.Mapper;
+
+/**
+ * 告警事件 Mapper
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Mapper
+public interface AlarmEventMapper extends BaseMapper<AlarmEvent> {
+}

+ 15 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/mapper/indicator/IndicatorResultMapper.java

@@ -0,0 +1,15 @@
+package com.hdkj.lt.mq.mapper.indicator;
+
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import com.hdkj.lt.mq.entity.indicator.IndicatorResult;
+import org.apache.ibatis.annotations.Mapper;
+
+/**
+ * 指标计算结果 Mapper
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Mapper
+public interface IndicatorResultMapper extends BaseMapper<IndicatorResult> {
+}

+ 28 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/AlarmEvaluateService.java

@@ -0,0 +1,28 @@
+package com.hdkj.lt.mq.service.indicator;
+
+import com.hdkj.lt.mq.entity.indicator.AlarmEvent;
+
+import java.util.List;
+
+/**
+ * 告警评估服务
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+public interface AlarmEvaluateService {
+
+    /**
+     * 评估告警事件
+     * <p>
+     * 解析算法返回的 rawResponse,从 over_limit_details 提取超限馈线,
+     * 从 feeder_results 提取馈线级指标值(负载率/电压),构建告警事件列表。
+     *
+     * @param indicatorResultId 指标结果ID
+     * @param batchId           批次ID
+     * @param countyId          区县ID
+     * @param rawResponseJson   算法返回的原始JSON(t_indicator_result.raw_response)
+     * @return 告警事件列表(按馈线粒度),无告警返回空列表
+     */
+    List<AlarmEvent> evaluateAlarms(Long indicatorResultId, String batchId, Integer countyId, String rawResponseJson);
+}

+ 31 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/IndicatorResultService.java

@@ -0,0 +1,31 @@
+package com.hdkj.lt.mq.service.indicator;
+
+import com.hdkj.lt.mq.entity.indicator.IndicatorResult;
+
+/**
+ * 指标计算结果服务
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+public interface IndicatorResultService {
+
+    /**
+     * 保存指标计算结果(先删后插,幂等)
+     *
+     * @param rawResponseJson 算法返回的原始JSON
+     * @param batchId         批次ID
+     * @param countyId        区县ID
+     * @return 入库后的指标结果
+     */
+    IndicatorResult saveResult(String rawResponseJson, String batchId, Integer countyId);
+
+    /**
+     * 按批次+区县查询指标结果
+     *
+     * @param batchId  批次ID
+     * @param countyId 区县ID
+     * @return 指标结果,不存在返回 null
+     */
+    IndicatorResult getByBatchAndCounty(String batchId, Integer countyId);
+}

+ 249 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/impl/AlarmEvaluateServiceImpl.java

@@ -0,0 +1,249 @@
+package com.hdkj.lt.mq.service.indicator.impl;
+
+import com.alibaba.fastjson.JSON;
+import com.alibaba.fastjson.JSONArray;
+import com.alibaba.fastjson.JSONObject;
+import com.hdkj.lt.mq.entity.indicator.AlarmEvent;
+import com.hdkj.lt.mq.service.indicator.AlarmEvaluateService;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Service;
+
+import java.math.BigDecimal;
+import java.time.LocalDate;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 告警评估服务实现
+ * <p>
+ * 解析算法返回的 rawResponse,分两步:
+ * <ol>
+ *   <li>解析 over_limit_details + feeder_results,提取超限馈线和指标值</li>
+ *   <li>调用 shouldAlarm() 自定义告警标志判断,满足条件才生成告警事件</li>
+ * </ol>
+ * <p>
+ * alarm_types 值域:
+ * <ul>
+ *   <li>load_rate     → 线路重过载</li>
+ *   <li>voltage_over  → 电压越限</li>
+ *   <li>voltage_under → 电压越限</li>
+ * </ul>
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Slf4j
+@Service
+public class AlarmEvaluateServiceImpl implements AlarmEvaluateService {
+
+    @Override
+    public List<AlarmEvent> evaluateAlarms(Long indicatorResultId, String batchId, Integer countyId, String rawResponseJson) {
+        log.info("评估告警 indicatorResultId={}, batchId={}, countyId={}, rawResponse={}",
+                indicatorResultId, batchId, countyId, rawResponseJson == null ? "null" : "provided");
+
+        List<AlarmEvent> alarmEvents = new ArrayList<AlarmEvent>();
+        if (StringUtils.isBlank(rawResponseJson)) {
+            log.warn("rawResponse 为空,跳过告警评估 indicatorResultId={}", indicatorResultId);
+            return alarmEvents;
+        }
+
+        JSONObject root = JSON.parseObject(rawResponseJson);
+        JSONObject calculationResult = root.getJSONObject("calculation_result");
+        if (calculationResult == null) {
+            log.warn("calculation_result 为空 indicatorResultId={}", indicatorResultId);
+            return alarmEvents;
+        }
+
+        JSONObject countyResult = calculationResult.getJSONObject("county_result");
+        if (countyResult == null) {
+            log.warn("county_result 为空 indicatorResultId={}", indicatorResultId);
+            return alarmEvents;
+        }
+
+        // ===== 第一步:解析超限数据 =====
+        // 从 over_limit_details 提取三类超限馈线
+        Map<String, List<String>> feederAlarmTypes = parseOverLimitDetails(countyResult);
+        if (feederAlarmTypes.isEmpty()) {
+            log.info("无超限馈线 indicatorResultId={}, batchId={}, countyId={}", indicatorResultId, batchId, countyId);
+            return alarmEvents;
+        }
+
+        // ===== 第二步:自定义告警标志判断 =====
+        if (!shouldAlarm(feederAlarmTypes, countyResult, root)) {
+            log.info("自定义告警判断未通过,不触发告警 indicatorResultId={}, batchId={}, countyId={}",
+                    indicatorResultId, batchId, countyId);
+            return alarmEvents;
+        }
+
+        // ===== 第三步:构建告警事件 =====
+        // 构建 feeder_results 索引(feeder_id → 指标值)
+        Map<String, JSONObject> feederResultMap = buildFeederResultMap(calculationResult);
+
+        // 取县域基本信息
+        JSONObject countyBasicInfo = root.getJSONObject("county_basic_info");
+        String countySid = countyBasicInfo != null ? countyBasicInfo.getString("county_sid") : null;
+        String countyName = countyBasicInfo != null ? countyBasicInfo.getString("county_name") : null;
+
+        // 为每个超限馈线构建 AlarmEvent
+        for (Map.Entry<String, List<String>> entry : feederAlarmTypes.entrySet()) {
+            String feederId = entry.getKey();
+            List<String> types = entry.getValue();
+            String alarmTypes = String.join(",", types);
+
+            // 从 feeder_results 取该馈线的指标值
+            JSONObject feederResult = feederResultMap.get(feederId);
+            BigDecimal maxLoadRatePct = null;
+            BigDecimal maxVoltagePu = null;
+            BigDecimal minVoltagePu = null;
+            if (feederResult != null) {
+                JSONObject loadRate = feederResult.getJSONObject("load_rate");
+                if (loadRate != null) {
+                    maxLoadRatePct = loadRate.getBigDecimal("max_load_rate_pct");
+                }
+                JSONObject voltagePu = feederResult.getJSONObject("voltage_pu");
+                if (voltagePu != null) {
+                    maxVoltagePu = voltagePu.getBigDecimal("max_voltage_pu");
+                    minVoltagePu = voltagePu.getBigDecimal("min_voltage_pu");
+                }
+            }
+
+            AlarmEvent alarmEvent = AlarmEvent.builder()
+                    .batchId(batchId)
+                    .countyId(countyId)
+                    .countySid(countySid)
+                    .countyName(countyName)
+                    .feederId(feederId)
+                    .alarmTypes(alarmTypes)
+                    .maxLoadRatePct(maxLoadRatePct)
+                    .maxVoltagePu(maxVoltagePu)
+                    .minVoltagePu(minVoltagePu)
+                    .indicatorResultId(indicatorResultId)
+                    .alarmDate(LocalDate.now())
+                    .build();
+
+            alarmEvents.add(alarmEvent);
+            log.info("告警评估-生成告警 feederId={}, alarmTypes={}, maxLoadRatePct={}, maxVoltagePu={}, minVoltagePu={}",
+                    feederId, alarmTypes, maxLoadRatePct, maxVoltagePu, minVoltagePu);
+        }
+
+        log.info("告警评估完成 indicatorResultId={}, batchId={}, countyId={}, alarmCount={}",
+                indicatorResultId, batchId, countyId, alarmEvents.size());
+        return alarmEvents;
+    }
+
+    // ============================================================
+    // 自定义告警标志判断(预留 customization point)
+    // ============================================================
+
+    /**
+     * 自定义告警标志判断
+     * <p>
+     * 自主判断指标数据是否满足告警条件。
+     * 满足返回 true → 继续触发下一步(生成告警事件)
+     * 不满足返回 false → 无事发生
+     * <p>
+     * 当前实现:算法返回 over_limit_details 有超限馈线即告警。
+     * 后续可在此自定义判断逻辑,例如:
+     * <ul>
+     *   <li>自定义阈值(不依赖算法的 flags,自行判断 max_load_rate_pct > 阈值)</li>
+     *   <li>白名单/黑名单(某些馈线不告警)</li>
+     *   <li>时间段过滤(某些时段不告警)</li>
+     *   <li>连续超限次数(N 次才告警)</li>
+     * </ul>
+     *
+     * @param feederAlarmTypes 超限馈线及告警类型(feederId → [load_rate / voltage_over / voltage_under])
+     * @param countyResult     县域汇总指标(load_rate / voltage_pu / loss / alarm)
+     * @param root             算法完整返回 JSON(含 county_basic_info / calculation_result)
+     * @return true=满足告警条件,false=不告警
+     */
+    private boolean shouldAlarm(Map<String, List<String>> feederAlarmTypes, JSONObject countyResult, JSONObject root) {
+        // TODO: 自定义告警标志判断逻辑
+        // 当前:算法判定超限即告警
+        return !feederAlarmTypes.isEmpty();
+    }
+
+    // ============================================================
+    // 数据解析
+    // ============================================================
+
+    /**
+     * 从 county_result.alarm.over_limit_details 解析三类超限馈线
+     *
+     * @return feederId → 告警类型列表(去重)
+     */
+    private Map<String, List<String>> parseOverLimitDetails(JSONObject countyResult) {
+        Map<String, List<String>> feederAlarmTypes = new LinkedHashMap<String, List<String>>();
+
+        JSONObject alarm = countyResult.getJSONObject("alarm");
+        if (alarm == null) {
+            return feederAlarmTypes;
+        }
+        JSONObject overLimitDetails = alarm.getJSONObject("over_limit_details");
+        if (overLimitDetails == null) {
+            return feederAlarmTypes;
+        }
+
+        List<String> loadRateFeeders = parseFeederList(overLimitDetails, "load_rate_over_limit");
+        List<String> voltageOverFeeders = parseFeederList(overLimitDetails, "voltage_over_upper_limit");
+        List<String> voltageUnderFeeders = parseFeederList(overLimitDetails, "voltage_under_lower_limit");
+
+        addAlarmTypes(feederAlarmTypes, loadRateFeeders, "load_rate");
+        addAlarmTypes(feederAlarmTypes, voltageOverFeeders, "voltage_over");
+        addAlarmTypes(feederAlarmTypes, voltageUnderFeeders, "voltage_under");
+
+        return feederAlarmTypes;
+    }
+
+    /**
+     * 从 over_limit_details 中解析某类超限馈线ID列表
+     */
+    private List<String> parseFeederList(JSONObject overLimitDetails, String key) {
+        JSONArray arr = overLimitDetails.getJSONArray(key);
+        if (arr == null || arr.isEmpty()) {
+            return new ArrayList<String>();
+        }
+        return arr.toJavaList(String.class);
+    }
+
+    /**
+     * 将馈线ID和告警类型加入 map
+     */
+    private void addAlarmTypes(Map<String, List<String>> feederAlarmTypes, List<String> feederIds, String alarmType) {
+        for (String feederId : feederIds) {
+            if (StringUtils.isBlank(feederId)) {
+                continue;
+            }
+            List<String> types = feederAlarmTypes.get(feederId);
+            if (types == null) {
+                types = new ArrayList<String>();
+                feederAlarmTypes.put(feederId, types);
+            }
+            if (!types.contains(alarmType)) {
+                types.add(alarmType);
+            }
+        }
+    }
+
+    /**
+     * 构建 feeder_results 索引:feeder_id → feederResult JSON 对象
+     */
+    private Map<String, JSONObject> buildFeederResultMap(JSONObject calculationResult) {
+        Map<String, JSONObject> map = new HashMap<String, JSONObject>();
+        JSONArray feederResults = calculationResult.getJSONArray("feeder_results");
+        if (feederResults == null || feederResults.isEmpty()) {
+            return map;
+        }
+        for (int i = 0; i < feederResults.size(); i++) {
+            JSONObject feeder = feederResults.getJSONObject(i);
+            String feederId = feeder.getString("feeder_id");
+            if (StringUtils.isNotBlank(feederId)) {
+                map.put(feederId, feeder);
+            }
+        }
+        return map;
+    }
+}

+ 175 - 0
services/load-transfer-mq/src/main/java/com/hdkj/lt/mq/service/indicator/impl/IndicatorResultServiceImpl.java

@@ -0,0 +1,175 @@
+package com.hdkj.lt.mq.service.indicator.impl;
+
+import com.alibaba.fastjson.JSON;
+import com.alibaba.fastjson.JSONArray;
+import com.alibaba.fastjson.JSONObject;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.hdkj.lt.mq.entity.indicator.IndicatorResult;
+import com.hdkj.lt.mq.mapper.indicator.IndicatorResultMapper;
+import com.hdkj.lt.mq.service.indicator.IndicatorResultService;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+
+import java.math.BigDecimal;
+import java.time.LocalDateTime;
+import java.util.List;
+
+/**
+ * 指标计算结果服务实现
+ *
+ * @author hermes
+ * @since 2026-07-17
+ */
+@Slf4j
+@Service
+@RequiredArgsConstructor
+public class IndicatorResultServiceImpl implements IndicatorResultService {
+
+    private final IndicatorResultMapper indicatorResultMapper;
+
+    @Override
+    public IndicatorResult saveResult(String rawResponseJson, String batchId, Integer countyId) {
+        log.info("保存指标计算结果 batchId={}, countyId={}", batchId, countyId);
+        JSONObject root = JSON.parseObject(rawResponseJson, JSONObject.class);
+
+        // county_basic_info
+        JSONObject countyBasicInfo = root.getJSONObject("county_basic_info");
+        String countySid = null;
+        String countyName = null;
+        if (countyBasicInfo != null) {
+            countySid = countyBasicInfo.getString("county_sid");
+            countyName = countyBasicInfo.getString("county_name");
+        }
+
+        // calculation_result
+        JSONObject calculationResult = root.getJSONObject("calculation_result");
+        JSONObject countyResult = calculationResult != null ? calculationResult.getJSONObject("county_result") : null;
+
+        // load_rate
+        JSONObject loadRate = countyResult != null ? countyResult.getJSONObject("load_rate") : null;
+        BigDecimal maxLoadRatePct = getBigDecimal(loadRate, "max_load_rate_pct");
+        BigDecimal minLoadRatePct = getBigDecimal(loadRate, "min_load_rate_pct");
+        BigDecimal avgLoadRatePct = getBigDecimal(loadRate, "avg_load_rate_pct");
+        BigDecimal alarmThresholdPct = getBigDecimal(loadRate, "alarm_threshold_pct");
+
+        // voltage_pu
+        JSONObject voltagePu = countyResult != null ? countyResult.getJSONObject("voltage_pu") : null;
+        BigDecimal maxVoltagePu = getBigDecimal(voltagePu, "max_voltage_pu");
+        BigDecimal minVoltagePu = getBigDecimal(voltagePu, "min_voltage_pu");
+        BigDecimal avgVoltagePu = getBigDecimal(voltagePu, "avg_voltage_pu");
+
+        // loss
+        JSONObject loss = countyResult != null ? countyResult.getJSONObject("loss") : null;
+        BigDecimal totalLossRatePct = getBigDecimal(loss, "total_loss_rate_pct");
+        BigDecimal transformerLossRatePct = getBigDecimal(loss, "transformer_loss_rate_pct");
+        BigDecimal lineLossRatePct = getBigDecimal(loss, "line_loss_rate_pct");
+
+        // alarm
+        JSONObject alarm = countyResult != null ? countyResult.getJSONObject("alarm") : null;
+        JSONObject alarmFlags = alarm != null ? alarm.getJSONObject("flags") : null;
+        Integer loadRateOverLimit = getInteger(alarmFlags, "load_rate_over_limit");
+        Integer voltageOverUpperLimit = getInteger(alarmFlags, "voltage_over_upper_limit");
+        Integer voltageUnderLowerLimit = getInteger(alarmFlags, "voltage_under_lower_limit");
+        Integer overLimitCount = getInteger(alarm, "over_limit_count");
+        String overLimitFeederIds = joinList(alarm, "over_limit_feeder_ids");
+
+        // reconfiguration_feeders.algorithm_trigger
+        JSONObject reconfigurationFeeders = calculationResult != null ? calculationResult.getJSONObject("reconfiguration_feeders") : null;
+        JSONObject algorithmTrigger = reconfigurationFeeders != null ? reconfigurationFeeders.getJSONObject("algorithm_trigger") : null;
+        Integer reconfigurationTriggered = getInteger(algorithmTrigger, "reconfiguration_triggered");
+        String triggerType = algorithmTrigger != null ? algorithmTrigger.getString("trigger_type") : null;
+        String problemFeederIds = joinList(algorithmTrigger, "problem_feeder_ids");
+
+        // 顶层字段
+        String status = root.getString("status");
+        String computeResult = root.getString("compute_result");
+        String generatedAt = root.getString("generated_at");
+
+        // 先删后插(幂等)
+        LambdaQueryWrapper<IndicatorResult> deleteWrapper = new LambdaQueryWrapper<IndicatorResult>()
+                .eq(IndicatorResult::getBatchId, batchId)
+                .eq(IndicatorResult::getCountyId, countyId);
+        indicatorResultMapper.delete(deleteWrapper);
+
+        IndicatorResult result = IndicatorResult.builder()
+                .batchId(batchId)
+                .countyId(countyId)
+                .countySid(countySid)
+                .countyName(countyName)
+                .status(status)
+                .computeResult(computeResult)
+                .generatedAt(generatedAt)
+                .maxLoadRatePct(maxLoadRatePct)
+                .minLoadRatePct(minLoadRatePct)
+                .avgLoadRatePct(avgLoadRatePct)
+                .alarmThresholdPct(alarmThresholdPct)
+                .maxVoltagePu(maxVoltagePu)
+                .minVoltagePu(minVoltagePu)
+                .avgVoltagePu(avgVoltagePu)
+                .totalLossRatePct(totalLossRatePct)
+                .transformerLossRatePct(transformerLossRatePct)
+                .lineLossRatePct(lineLossRatePct)
+                .loadRateOverLimit(loadRateOverLimit)
+                .voltageOverUpperLimit(voltageOverUpperLimit)
+                .voltageUnderLowerLimit(voltageUnderLowerLimit)
+                .overLimitCount(overLimitCount)
+                .overLimitFeederIds(overLimitFeederIds)
+                .reconfigurationTriggered(reconfigurationTriggered)
+                .triggerType(triggerType)
+                .problemFeederIds(problemFeederIds)
+                .rawResponse(rawResponseJson)
+                .createTime(LocalDateTime.now())
+                .build();
+
+        indicatorResultMapper.insert(result);
+        log.info("保存指标计算结果完成 batchId={}, countyId={}, id={}", batchId, countyId, result.getId());
+        return result;
+    }
+
+    @Override
+    public IndicatorResult getByBatchAndCounty(String batchId, Integer countyId) {
+        LambdaQueryWrapper<IndicatorResult> wrapper = new LambdaQueryWrapper<IndicatorResult>()
+                .eq(IndicatorResult::getBatchId, batchId)
+                .eq(IndicatorResult::getCountyId, countyId);
+        return indicatorResultMapper.selectOne(wrapper);
+    }
+
+    /**
+     * 从 JSONObject 安全获取 BigDecimal
+     */
+    private BigDecimal getBigDecimal(JSONObject obj, String key) {
+        if (obj == null || obj.get(key) == null) {
+            return null;
+        }
+        return obj.getBigDecimal(key);
+    }
+
+    /**
+     * 从 JSONObject 安全获取 Integer
+     */
+    private Integer getInteger(JSONObject obj, String key) {
+        if (obj == null || obj.get(key) == null) {
+            return null;
+        }
+        return obj.getInteger(key);
+    }
+
+    /**
+     * 从 JSONObject 获取 List<String> 并用逗号拼接
+     */
+    private String joinList(JSONObject obj, String key) {
+        if (obj == null) {
+            return null;
+        }
+        JSONArray arr = obj.getJSONArray(key);
+        if (arr == null || arr.isEmpty()) {
+            return null;
+        }
+        List<String> list = arr.toJavaList(String.class);
+        if (list.isEmpty()) {
+            return null;
+        }
+        return String.join(",", list);
+    }
+}

+ 33 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/controller/algorithm/AlgorithmController.java

@@ -0,0 +1,33 @@
+package com.hdkj.lt.si.controller.algorithm;
+
+import com.hdkj.hussar.ApiResponse;
+import com.hdkj.lt.core.mvc.BaseController;
+import com.hdkj.lt.si.modle.algorithm.IndicatorStatusResponse;
+import com.hdkj.lt.si.service.algorithm.AlgorithmService;
+import lombok.AllArgsConstructor;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+/**
+ * 算法服务接口
+ *
+ * @Description 对外暴露指标计算接口,前端每15分钟调用一次
+ */
+@RestController
+@AllArgsConstructor
+@RequestMapping("/algorithm")
+public class AlgorithmController extends BaseController {
+
+    private final AlgorithmService algorithmService;
+
+    /**
+     * 查询指标计算结果
+     *
+     * @return 指标计算结果
+     */
+    @PostMapping(value = "/indicators")
+    public ApiResponse<IndicatorStatusResponse> queryIndicators() {
+        return data(algorithmService.queryIndicators());
+    }
+}

+ 25 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/Alarm.java

@@ -0,0 +1,25 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 告警信息
+ */
+@Data
+public class Alarm {
+
+    private Flags flags;
+    private Integer overLimitCount;
+    private List<String> overLimitFeederIds;
+    private Map<String, List<String>> overLimitDetails;
+
+    @Data
+    public static class Flags {
+        private Integer loadRateOverLimit;
+        private Integer voltageOverUpperLimit;
+        private Integer voltageUnderLowerLimit;
+    }
+}

+ 19 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/AlgorithmTrigger.java

@@ -0,0 +1,19 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+import java.util.List;
+
+/**
+ * 重构算法触发信息
+ */
+@Data
+public class AlgorithmTrigger {
+
+    private Integer triggered;
+    private String triggerType;
+    private String reason;
+    private Double thresholdPct;
+    private List<String> problemFeederIds;
+    private List<String> triggerFeederIds;
+}

+ 15 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/CountyResult.java

@@ -0,0 +1,15 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+/**
+ * 县域指标结果
+ */
+@Data
+public class CountyResult {
+
+    private LoadRate loadRate;
+    private VoltagePu voltagePu;
+    private Loss loss;
+    private Alarm alarm;
+}

+ 17 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/FeederResult.java

@@ -0,0 +1,17 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+/**
+ * 馈线指标结果
+ */
+@Data
+public class FeederResult {
+
+    private String feederId;
+    private String feederName;
+    private LoadRate loadRate;
+    private VoltagePu voltagePu;
+    private Loss loss;
+    private Alarm alarm;
+}

+ 44 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/IndicatorStatusResponse.java

@@ -0,0 +1,44 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+import java.util.List;
+
+/**
+ * 算法服务指标计算响应
+ *
+ * @Description 观音阁算法服务 /api/indicators 返回结构
+ */
+@Data
+public class IndicatorStatusResponse {
+
+    private String status;
+
+    private String responseType;
+
+    private String computeResult;
+
+    private String generatedAt;
+
+    private CountyBasicInfo countyBasicInfo;
+
+    private CalculationResult calculationResult;
+
+    @Data
+    public static class CountyBasicInfo {
+        private Integer countyId;
+        private String countySid;
+        private String countyName;
+        private List<String> feederIds;
+        private Boolean usedDefaultGroup;
+        private String groupSource;
+        private String countySource;
+    }
+
+    @Data
+    public static class CalculationResult {
+        private CountyResult countyResult;
+        private List<FeederResult> feederResults;
+        private ReconfigurationFeeders reconfigurationFeeders;
+    }
+}

+ 15 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/LoadRate.java

@@ -0,0 +1,15 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+/**
+ * 负载率指标
+ */
+@Data
+public class LoadRate {
+
+    private Double maxLoadRatePct;
+    private Double minLoadRatePct;
+    private Double avgLoadRatePct;
+    private Double alarmThresholdPct;
+}

+ 15 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/Loss.java

@@ -0,0 +1,15 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+/**
+ * 损耗指标
+ */
+@Data
+public class Loss {
+
+    private Double totalLossRatePct;
+    private Double transformerLossRatePct;
+    private Double lineLossRatePct;
+    private String unit;
+}

+ 18 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/ReconfigurationFeeders.java

@@ -0,0 +1,18 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+import java.util.List;
+
+/**
+ * 重构馈线范围及触发标志
+ */
+@Data
+public class ReconfigurationFeeders {
+
+    private Integer groupId;
+    private String groupSid;
+    private String groupName;
+    private List<String> feederIds;
+    private AlgorithmTrigger algorithmTrigger;
+}

+ 14 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/modle/algorithm/VoltagePu.java

@@ -0,0 +1,14 @@
+package com.hdkj.lt.si.modle.algorithm;
+
+import lombok.Data;
+
+/**
+ * 电压标幺值指标
+ */
+@Data
+public class VoltagePu {
+
+    private Double maxVoltagePu;
+    private Double minVoltagePu;
+    private Double avgVoltagePu;
+}

+ 16 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/algorithm/AlgorithmService.java

@@ -0,0 +1,16 @@
+package com.hdkj.lt.si.service.algorithm;
+
+import com.hdkj.lt.si.modle.algorithm.IndicatorStatusResponse;
+
+/**
+ * 算法服务业务接口
+ */
+public interface AlgorithmService {
+
+    /**
+     * 查询指标计算结果
+     *
+     * @return 指标计算结果
+     */
+    IndicatorStatusResponse queryIndicators();
+}

+ 27 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/algorithm/impl/AlgorithmServiceImpl.java

@@ -0,0 +1,27 @@
+package com.hdkj.lt.si.service.algorithm.impl;
+
+import com.hdkj.lt.si.modle.algorithm.IndicatorStatusResponse;
+import com.hdkj.lt.si.service.algorithm.AlgorithmService;
+import com.hdkj.lt.si.service.remote.impl.AlgorithmRemoteCallServiceImpl;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+
+import javax.annotation.Resource;
+
+/**
+ * 算法服务业务实现
+ */
+@Slf4j
+@Service
+@RequiredArgsConstructor
+public class AlgorithmServiceImpl implements AlgorithmService {
+
+    @Resource(type = AlgorithmRemoteCallServiceImpl.class)
+    private AlgorithmRemoteCallServiceImpl algorithmRemoteCallService;
+
+    @Override
+    public IndicatorStatusResponse queryIndicators() {
+        return algorithmRemoteCallService.queryIndicators();
+    }
+}

+ 74 - 0
services/load-transfer-si/src/main/java/com/hdkj/lt/si/service/remote/impl/AlgorithmRemoteCallServiceImpl.java

@@ -0,0 +1,74 @@
+package com.hdkj.lt.si.service.remote.impl;
+
+import com.alibaba.fastjson.JSON;
+import com.hdkj.lt.base.exception.ExceptionCast;
+import com.hdkj.lt.core.rest.config.HttpClientConfig;
+import com.hdkj.lt.si.modle.algorithm.IndicatorStatusResponse;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.codec.Charsets;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.MediaType;
+import org.springframework.http.ResponseEntity;
+import org.springframework.stereotype.Service;
+import org.springframework.web.client.RestTemplate;
+
+import javax.annotation.Resource;
+import java.util.Objects;
+
+/**
+ * 算法服务远程调用
+ *
+ * @Description 调用观音阁算法服务 /api/indicators 接口,获取指标计算结果。
+ * 算法服务直接返回业务数据(无 code/data 外壳),不继承 AbstractDefaultRemoteCallService。
+ */
+@Slf4j
+@Service
+public class AlgorithmRemoteCallServiceImpl {
+
+    @Resource(name = HttpClientConfig.CONN_POOL_REST_TEMPLATE)
+    private RestTemplate restTemplate;
+
+    @Value("${remote.algorithm.request-uri}")
+    private String requestUri;
+
+    @Value("${remote.algorithm.indicators-uri}")
+    private String indicatorsUri;
+
+    /**
+     * 调用指标计算接口
+     *
+     * @return 指标计算结果
+     */
+    public IndicatorStatusResponse queryIndicators() {
+        String requestUrl = requestUri.concat(indicatorsUri);
+        String requestBody = "{}";
+        try {
+            HttpHeaders headers = new HttpHeaders();
+            headers.set(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE);
+            headers.set(HttpHeaders.ACCEPT_CHARSET, Charsets.UTF_8.toString());
+            HttpEntity<String> httpEntity = new HttpEntity<>(requestBody, headers);
+
+            log.info("调用算法指标接口, requestUrl: {}", requestUrl);
+            ResponseEntity<IndicatorStatusResponse> responseEntity = restTemplate.postForEntity(
+                    requestUrl, httpEntity, IndicatorStatusResponse.class);
+
+            if (!responseEntity.getStatusCode().is2xxSuccessful() || Objects.isNull(responseEntity.getBody())) {
+                String errorMsg = String.format("调用算法指标接口失败, requestUrl: %s, response: %s",
+                        requestUrl, JSON.toJSONString(responseEntity));
+                log.error(errorMsg);
+                ExceptionCast.cast(errorMsg);
+            }
+
+            IndicatorStatusResponse result = responseEntity.getBody();
+            log.info("调用算法指标接口成功, requestUrl: {}, computeResult: {}",
+                    requestUrl, result.getComputeResult());
+            return result;
+        } catch (Exception e) {
+            log.error("调用算法指标接口异常, requestUrl: {}, errorMsg: {}", requestUrl, e.toString());
+            ExceptionCast.cast("调用算法指标接口失败: " + e.getMessage());
+            return null;
+        }
+    }
+}