liuaini 3 днів тому
батько
коміт
a1d16401ec

+ 129 - 0
common/common-se/src/main/java/com/hdkj/lt/bf/scheduler/SePartitionRotateScheduler.java

@@ -0,0 +1,129 @@
+package com.hdkj.lt.bf.scheduler;
+
+import com.google.common.collect.ImmutableMap;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.stereotype.Component;
+
+import java.time.LocalDate;
+import java.time.format.DateTimeFormatter;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 运行态明细分区滚动清理(公共模块)
+ * <p>
+ * 背景:fhzg_se_snapshot_detail(断面明细)/ fhzg_se_line_loss_detail(线损明细)
+ * 为按月 RANGE COLUMNS 分区表,需每月滚动:
+ * 1. REORGANIZE pMax 拆出下月分区(pMax 含跨月数据时不能直接 ADD,必须 REORGANIZE)
+ * 2. DROP 过期月分区(保留策略:断面 30 天 / 线损 7 天 → 最多保留近 2 个自然月)
+ * <p>
+ * 前置条件:表已完成分区迁移(scripts/migrate_se_detail_partition.sql)。
+ * 表未分区时安全跳过并告警,不抛异常。
+ * 定时调度由 pwfhzg-job 通过 sys_job 配置触发(每月1号凌晨2点)。
+ *
+ * @author lsl
+ * @since 2026-08-07
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class SePartitionRotateScheduler {
+
+    private static final DateTimeFormatter PART_FMT = DateTimeFormatter.ofPattern("yyyyMM");
+
+    /** 需要滚动的分区表:表名 → 保留月数(含当前月,即最多保留 N+1 个自然月) */
+    private static final ImmutableMap<String, Integer> PARTITION_TABLES = ImmutableMap.of(
+            "fhzg_se_snapshot_detail", 2,    // 30 天保留 → 当前月+上月
+            "fhzg_se_line_loss_detail", 2     // 7 天保留 → 当前月+上月
+    );
+
+    private final JdbcTemplate jdbcTemplate;
+
+    /**
+     * 每月执行一次:滚动分区(拆下月 + 删过期)
+     *
+     * @return 每张表的执行结果描述
+     */
+    public String rotateMonthly() {
+        LocalDate today = LocalDate.now();
+        String nextMonthPart = "p" + today.plusMonths(1).format(PART_FMT);
+        String nextMonthBound = today.plusMonths(2).withDayOfMonth(1).toString();   // 下月1日 → 下月分区的上界
+        String keepFromMonth = today.withDayOfMonth(1).minusMonths(1).toString();   // 保留 >= 上月1日
+
+        StringBuilder summary = new StringBuilder();
+        for (Map.Entry<String, Integer> entry : PARTITION_TABLES.entrySet()) {
+            String table = entry.getKey();
+            int keepMonths = entry.getValue();
+            try {
+                summary.append(rotateOneTable(table, nextMonthPart, nextMonthBound, keepFromMonth, keepMonths));
+            } catch (Exception e) {
+                log.error("[分区滚动] {} 执行异常", table, e);
+                summary.append(table).append(" 执行异常: ").append(e.getMessage()).append("; ");
+            }
+        }
+        log.info("[分区滚动] 月度执行完成: {}", summary);
+        return summary.toString();
+    }
+
+    private String rotateOneTable(String table, String nextMonthPart, String nextMonthBound,
+                                  String keepFromMonth, int keepMonths) {
+        // 1. 检查表是否已分区(防 migrate 未执行时误操作)
+        Integer partCount = jdbcTemplate.queryForObject(
+                "SELECT COUNT(*) FROM information_schema.PARTITIONS " +
+                        "WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND PARTITION_NAME IS NOT NULL",
+                Integer.class, table);
+        if (partCount == null || partCount == 0) {
+            String warn = table + " 未分区,跳过滚动(请先执行 migrate_se_detail_partition.sql); ";
+            log.warn("[分区滚动] {}", warn);
+            return warn;
+        }
+
+        // 2. 查当前分区列表(名称 + 上界)
+        List<Map<String, Object>> partitions = jdbcTemplate.queryForList(
+                "SELECT PARTITION_NAME, PARTITION_DESCRIPTION FROM information_schema.PARTITIONS " +
+                        "WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND PARTITION_NAME IS NOT NULL " +
+                        "ORDER BY PARTITION_ORDINAL_POSITION", table);
+
+        // 3. REORGANIZE pMax 拆下月分区(每次执行都先拆,幂等:若下月分区已存在则跳过)
+        boolean hasNextPart = partitions.stream()
+                .anyMatch(p -> nextMonthPart.equals(p.get("PARTITION_NAME")));
+        if (!hasNextPart) {
+            String reorganizeSql = "ALTER TABLE `" + table + "` REORGANIZE PARTITION pMax INTO (" +
+                    "PARTITION " + nextMonthPart + " VALUES LESS THAN ('" + nextMonthBound + "'), " +
+                    "PARTITION pMax VALUES LESS THAN (MAXVALUE))";
+            jdbcTemplate.execute(reorganizeSql);
+            log.info("[分区滚动] {} REORGANIZE 完成: 新增分区 {}", table, nextMonthPart);
+            partitions = jdbcTemplate.queryForList(
+                    "SELECT PARTITION_NAME, PARTITION_DESCRIPTION FROM information_schema.PARTITIONS " +
+                            "WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND PARTITION_NAME IS NOT NULL " +
+                            "ORDER BY PARTITION_ORDINAL_POSITION", table);
+        } else {
+            log.info("[分区滚动] {} 下月分区 {} 已存在,跳过 REORGANIZE", table, nextMonthPart);
+        }
+
+        // 4. DROP 过期分区(保留 keepMonths 个自然月)
+        StringBuilder dropped = new StringBuilder();
+        for (Map<String, Object> p : partitions) {
+            String name = (String) p.get("PARTITION_NAME");
+            Object desc = p.get("PARTITION_DESCRIPTION");
+            if (name == null || "pMax".equals(name) || name.equals(nextMonthPart)) continue;
+            // PARTITION_DESCRIPTION 是 'yyyy-MM-dd' 格式的字符串(RANGE COLUMNS)
+            String bound = desc != null ? desc.toString() : "";
+            // 只处理形如 p202608 的月分区(非 p 开头数字的跳过)
+            if (!bound.matches("'?\\d{4}-\\d{2}-\\d{2}'?.*")) continue;
+            String cleanBound = bound.replaceAll("'", "").substring(0, 10);
+            // 保留窗口:分区上界 > 上月1日 才保留;上界 <= 上月1日 → 整个分区已过期 → DROP
+            // 例:9-01 跑时 keepFromMonth=08-01,p202607(上界08-01)<=08-01 → 删;p202608(上界09-01)>08-01 → 保留
+            if (cleanBound.compareTo(keepFromMonth) > 0) continue;
+            jdbcTemplate.execute("ALTER TABLE `" + table + "` DROP PARTITION " + name);
+            dropped.append(name).append(" ");
+        }
+        String result = table + " 滚动完成(拆" + nextMonthPart + (dropped.length() > 0 ? ", 删 " + dropped.toString().trim() : ", 无过期分区") + "); ";
+        log.info("[分区滚动] {}", result);
+        return result;
+    }
+}

+ 9 - 0
common/pom.xml

@@ -45,6 +45,15 @@
             <groupId>com.baomidou</groupId>
             <artifactId>mybatis-plus-boot-starter</artifactId>
         </dependency>
+        <dependency>
+            <groupId>com.google.code.gson</groupId>
+            <artifactId>gson</artifactId>
+        </dependency>
+        <dependency>
+            <groupId>com.google.guava</groupId>
+            <artifactId>guava</artifactId>
+            <version>31.1-jre</version>
+        </dependency>
     </dependencies>
 
 </project>

+ 2 - 5
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/common/SeMonitorThreshold.java

@@ -8,7 +8,7 @@ import java.math.BigDecimal;
  * 馈线电压: ±7% (ratio 0.93~1.07)
  * 配变电压: -10% ~ +7% (ratio 0.90~1.07)
  * 用户0.4kV: ±10%一般, -20%~+15%严重 (ratio 0.80~1.10一般, 0.80~0.90/1.10~1.15严重)
- * 负载率: ≥80%重载, ≥100%过载
+ * 负载率: ≥75%重载, ≥100%过载
  *
  * @author lsl
  * @since 2026-07-23
@@ -34,12 +34,9 @@ public final class SeMonitorThreshold {
     public static final BigDecimal CONSUMER_SEVERE_V_UNDER = new BigDecimal("0.80");
 
     // === 负载率 ===
-    public static final BigDecimal LOAD_RATE_HEAVY = new BigDecimal("80");
+    public static final BigDecimal LOAD_RATE_HEAVY = new BigDecimal("75");
     public static final BigDecimal LOAD_RATE_OVERLOAD = new BigDecimal("100");
 
-    // === 电流重过载持续触发: 时间跨度 ≥ 15分钟 ===
-    public static final long CURRENT_TRIGGER_DURATION_MINUTES = 15L;
-
     // === 电压越限持续触发: 时间跨度 ≥ 60分钟 ===
     public static final long TRIGGER_DURATION_MINUTES = 60L;
 

+ 3 - 3
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/common/SeSnapshotConstants.java

@@ -18,8 +18,8 @@ public interface SeSnapshotConstants {
 
     /** 区县列表:maintOrg/countyId → [countySid, countyName] */
     List<Map.Entry<String, String[]>> COUNTIES = Arrays.asList(
-            new AbstractMap.SimpleEntry<>("BA6481B8F34143FDA7E1C35A76846911", new String[]{"洪湖", "国网洪湖市供电公司"}),
-            new AbstractMap.SimpleEntry<>("54C4BBBB8FE4428E96B6AE12247C4DD4", new String[]{"秭归", "国网秭归县供电公司"}),
-            new AbstractMap.SimpleEntry<>("60FDC66B3FC143DC9C39001215DDCC61", new String[]{"长阳", "国网长阳县供电公司"})
+            new AbstractMap.SimpleEntry<>("BA6481B8F34143FDA7E1C35A76846911", new String[]{"honghu", "国网洪湖市供电公司"}),
+            new AbstractMap.SimpleEntry<>("54C4BBBB8FE4428E96B6AE12247C4DD4", new String[]{"zigui", "国网秭归县供电公司"}),
+            new AbstractMap.SimpleEntry<>("60FDC66B3FC143DC9C39001215DDCC61", new String[]{"changyang", "国网长阳县供电公司"})
     );
 }

+ 51 - 58
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/impl/SeIndicatorServiceImpl.java

@@ -193,24 +193,26 @@ public class SeIndicatorServiceImpl implements SeIndicatorService {
      */
     private int[] countCurrentEvents(String id, Integer type,
                                       LocalDateTime start, LocalDateTime end) {
-        LambdaQueryWrapper<FhzgSeCurrentEvent> wrapper = new LambdaQueryWrapper<FhzgSeCurrentEvent>()
-                .ge(FhzgSeCurrentEvent::getFirstOverTime, start)
-                .lt(FhzgSeCurrentEvent::getFirstOverTime, end)
-                .in(FhzgSeCurrentEvent::getAlarmType, "current_heavy", "current_overload");
+        QueryWrapper<FhzgSeCurrentEvent> qw = new QueryWrapper<>();
+        qw.select("COALESCE(SUM(CASE WHEN alarm_type = 'current_heavy' THEN 1 ELSE 0 END), 0) AS heavy",
+                "COALESCE(SUM(CASE WHEN alarm_type = 'current_overload' THEN 1 ELSE 0 END), 0) AS overload")
+                .ge("first_over_time", start)
+                .lt("first_over_time", end);
         if (type == 2) {
-            wrapper.eq(FhzgSeCurrentEvent::getCountyId, id);
+            qw.eq("county_id", id);
         } else if (type == 3) {
-            wrapper.eq(FhzgSeCurrentEvent::getSubsId, id);
+            qw.eq("subs_id", id);
         } else if (type == 4) {
-            wrapper.eq(FhzgSeCurrentEvent::getFeederId, id);
+            qw.eq("feeder_id", id);
         }
 
-        List<FhzgSeCurrentEvent> events = currentEventMapper.selectList(wrapper);
-        int heavy = 0, overload = 0;
-        for (FhzgSeCurrentEvent e : events) {
-            if ("current_heavy".equals(e.getAlarmType())) heavy++;
-            else if ("current_overload".equals(e.getAlarmType())) overload++;
+        List<Map<String, Object>> rows = currentEventMapper.selectMaps(qw);
+        if (rows == null || rows.isEmpty() || rows.get(0) == null) {
+            return new int[]{0, 0};
         }
+        Map<String, Object> row = rows.get(0);
+        int heavy = ((Number) row.get("heavy")).intValue();
+        int overload = ((Number) row.get("overload")).intValue();
         return new int[]{heavy, overload};
     }
 
@@ -401,28 +403,41 @@ public class SeIndicatorServiceImpl implements SeIndicatorService {
      * 今日供电能力:负载率从断面明细聚合,重过载次数从事件表统计
      */
     private SeCapacityDashboardVO queryCapacityToday(String id, Integer type) {
-        LocalDate today = LocalDate.now();
-        LocalDateTime startTime = today.atStartOfDay();
+        LocalDateTime startTime = LocalDate.now().atStartOfDay();
         LocalDateTime endTime = LocalDateTime.now();
 
-        // 1. 负载率:从断面明细聚合(保持馈线96点时序)
-        LambdaQueryWrapper<FhzgSeSnapshotDetail> snapshotWrapper = new LambdaQueryWrapper<FhzgSeSnapshotDetail>()
-                .eq(FhzgSeSnapshotDetail::getDeviceType, "feeder")
-                .isNotNull(FhzgSeSnapshotDetail::getLoadRate)
-                .ge(FhzgSeSnapshotDetail::getSnapshotTime, startTime)
-                .le(FhzgSeSnapshotDetail::getSnapshotTime, endTime);
-        applySnapshotFilter(snapshotWrapper, id, type);
-        List<FhzgSeSnapshotDetail> snapshots = snapshotDetailMapper.selectList(snapshotWrapper);
-
-        SeCapacityDashboardVO vo;
-        if (snapshots.isEmpty()) {
-            vo = SeCapacityDashboardVO.builder().build();
-        } else {
-            vo = aggregateSnapshots(snapshots);
-            // 馈线 → 附加96点时序
-            if (type == 4) {
-                vo.setLoadRateSeries(buildFeederTimeSeries(snapshots));
-            }
+        // 1. 负载率:SQL 聚合一步到位(MAX/AVG),不再全量拉明细
+        QueryWrapper<FhzgSeSnapshotDetail> aggWrapper = new QueryWrapper<>();
+        aggWrapper.select("MAX(load_rate) AS max_load_rate",
+                "AVG(load_rate) AS avg_load_rate")
+                .eq("device_type", "feeder")
+                .isNotNull("load_rate")
+                .ge("snapshot_time", startTime)
+                .le("snapshot_time", endTime);
+        applySnapshotFilter(aggWrapper, id, type);
+        List<Map<String, Object>> aggRows = snapshotDetailMapper.selectMaps(aggWrapper);
+        SeCapacityDashboardVO vo = SeCapacityDashboardVO.builder().build();
+        if (aggRows != null && !aggRows.isEmpty() && aggRows.get(0) != null
+                && aggRows.get(0).get("max_load_rate") != null) {
+            Map<String, Object> row = aggRows.get(0);
+            vo.setMaxLoadRate(new BigDecimal(row.get("max_load_rate").toString()));
+            vo.setAvgLoadRate(new BigDecimal(row.get("avg_load_rate").toString())
+                    .setScale(2, RoundingMode.HALF_UP));
+        }
+
+        // 馈线 → 附加96点时序(仅 type=4 需要,限制列拉取)
+        if (type == 4) {
+            LambdaQueryWrapper<FhzgSeSnapshotDetail> seriesWrapper = new LambdaQueryWrapper<FhzgSeSnapshotDetail>()
+                    .select(FhzgSeSnapshotDetail::getSnapshotTime, FhzgSeSnapshotDetail::getLoadRate,
+                            FhzgSeSnapshotDetail::getCurrentValue, FhzgSeSnapshotDetail::getVoltageValue,
+                            FhzgSeSnapshotDetail::getActivePower, FhzgSeSnapshotDetail::getReactivePower)
+                    .eq(FhzgSeSnapshotDetail::getDeviceType, "feeder")
+                    .eq(FhzgSeSnapshotDetail::getFeederId, id)
+                    .isNotNull(FhzgSeSnapshotDetail::getLoadRate)
+                    .ge(FhzgSeSnapshotDetail::getSnapshotTime, startTime)
+                    .le(FhzgSeSnapshotDetail::getSnapshotTime, endTime);
+            List<FhzgSeSnapshotDetail> snapshots = snapshotDetailMapper.selectList(seriesWrapper);
+            vo.setLoadRateSeries(buildFeederTimeSeries(snapshots));
         }
 
         // 2. 重过载事件次数:从事件表统计
@@ -552,28 +567,6 @@ public class SeIndicatorServiceImpl implements SeIndicatorService {
         return vo;
     }
 
-    /**
-     * 内存聚合断面列表 → VO(today用,只算负载率,事件次数由 queryCapacityToday 从事件表统计)
-     */
-    private SeCapacityDashboardVO aggregateSnapshots(List<FhzgSeSnapshotDetail> snapshots) {
-        BigDecimal maxLoadRate = snapshots.stream()
-                .map(FhzgSeSnapshotDetail::getLoadRate)
-                .filter(Objects::nonNull)
-                .reduce(BigDecimal::max)
-                .orElse(BigDecimal.ZERO);
-
-        BigDecimal avgLoadRate = snapshots.stream()
-                .map(FhzgSeSnapshotDetail::getLoadRate)
-                .filter(Objects::nonNull)
-                .reduce(BigDecimal.ZERO, BigDecimal::add)
-                .divide(BigDecimal.valueOf(snapshots.size()), 2, RoundingMode.HALF_UP);
-
-        return SeCapacityDashboardVO.builder()
-                .maxLoadRate(maxLoadRate)
-                .avgLoadRate(avgLoadRate)
-                .build();
-    }
-
     /**
      * 构建馈线今日96点时序(负载率+电流+电压+有功+无功)
      */
@@ -606,13 +599,13 @@ public class SeIndicatorServiceImpl implements SeIndicatorService {
         }
     }
 
-    private void applySnapshotFilter(LambdaQueryWrapper<FhzgSeSnapshotDetail> wrapper, String id, Integer type) {
+    private void applySnapshotFilter(QueryWrapper<FhzgSeSnapshotDetail> wrapper, String id, Integer type) {
         if (type == 2) {
-            wrapper.eq(FhzgSeSnapshotDetail::getCountyId, id);
+            wrapper.eq("county_id", id);
         } else if (type == 3) {
-            wrapper.eq(FhzgSeSnapshotDetail::getSubsId, id);
+            wrapper.eq("subs_id", id);
         } else if (type == 4) {
-            wrapper.eq(FhzgSeSnapshotDetail::getFeederId, id);
+            wrapper.eq("feeder_id", id);
         }
     }
 

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

@@ -236,9 +236,9 @@ public class SeSnapshotServiceImpl implements SeSnapshotService {
             if (isOverLimit) {
                 if (state.getFirstOverTime() == null) state.setFirstOverTime(snapTime);
                 if (state.getFirstOverTime() != null
-                        && Duration.between(state.getFirstOverTime(), snapTime).toMinutes() >= SeMonitorThreshold.CURRENT_TRIGGER_DURATION_MINUTES
                         && state.getAlarmTriggered() == 0) {
-                    // 触发时往回看持续窗口:有任意断面≥100%算过载,否则重载
+                    // 一监测到越限立即生成事件(无持续时长判定)
+                    // 触发时往回看当前断面:≥100%算过载,否则重载
                     boolean hasOverload = hasOverloadSnapshotInWindow(feederId, state.getFirstOverTime(), snapTime);
                     String alarmType = hasOverload ? "current_overload" : "current_heavy";
                     String countyName = lookupCountyName(countyId, feederMap);

+ 0 - 1
services/load-transfer-bf/src/main/java/com/hdkj/lt/bf/service/staticgrid/impl/StaticGridDetailServiceImpl.java

@@ -52,7 +52,6 @@ public class StaticGridDetailServiceImpl
 
     private final StaticIndicatorMetricMapper metricMapper;
     private final StaticProblemDetailMapper problemDetailMapper;
-    private final RedisRepository redisRepository;
 
     @Override
     public StaticGridDetailResultVO queryDetail(StaticGridDetailQueryDTO query) {

+ 38 - 0
services/ruoyi-job/src/main/java/com/hdkj/lt/job/task/SeDataCleanupTask.java

@@ -0,0 +1,38 @@
+package com.hdkj.lt.job.task;
+
+import com.hdkj.lt.bf.scheduler.SePartitionRotateScheduler;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Component;
+
+/**
+ * 运行态明细分区滚动清理定时任务(每月1号凌晨2点)
+ * <p>
+ * sys_job 配置:invokeTarget = seDataCleanup.noParams(),cron = 0 0 2 1 * ?
+ * <p>
+ * 职责(委派 common-se SePartitionRotateScheduler):
+ * 1. fhzg_se_snapshot_detail:REORGANIZE pMax 拆下月分区 + DROP 30 天前过期分区
+ * 2. fhzg_se_line_loss_detail:REORGANIZE pMax 拆下月分区 + DROP 7 天前过期分区
+ * <p>
+ * 前置:表已完成分区迁移(scripts/migrate_se_detail_partition.sql);未分区时安全跳过并告警。
+ *
+ * @author lsl
+ * @since 2026-08-07
+ */
+@Slf4j
+@Component("seDataCleanup")
+@RequiredArgsConstructor
+public class SeDataCleanupTask {
+
+    private final SePartitionRotateScheduler sePartitionRotateScheduler;
+
+    public void noParams() {
+        log.info("[分区滚动] 开始执行月度分区清理");
+        try {
+            String result = sePartitionRotateScheduler.rotateMonthly();
+            log.info("[分区滚动] 执行完成: {}", result);
+        } catch (Exception e) {
+            log.error("[分区滚动] 执行失败", e);
+        }
+    }
+}