From 8290b813f1d36c48e6c110e77ab3fdc371614346 Mon Sep 17 00:00:00 2001 From: bot_pm Date: Sun, 14 Jun 2026 13:27:31 +0800 Subject: [PATCH] =?UTF-8?q?Phase=202=20#10=20#11=20#12=20#13=20#14=20#15:?= =?UTF-8?q?=20=E4=BE=9B=E6=B0=B4=E7=94=9F=E4=BA=A7=E7=AE=A1=E7=90=86?= =?UTF-8?q?=E5=B9=B3=E5=8F=B0=20+=20=E5=B7=A1=E6=A3=80=E7=AE=A1=E7=90=86?= =?UTF-8?q?=E7=B3=BB=E7=BB=9F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #10 总览+在线监测: - DashboardService: 今日进出水量/设备概况/能耗药耗/实时监测列表(多维筛选) - VideoService: 视频监控点位+AI人员闯入检测(YOLO mock) #11 水质管控+报警: - WaterQualityService: 全工艺药剂投加监控(混凝/沉淀/过滤/消毒) + 水质台账 - AlertEngine: 报警规则检测/去重/确认/派单/分级(info/warning/critical/emergency) #12 调度工作台+调度业务: - DispatchService: 值班管理(开始/结束/交接) + 指令创建/下发/跟踪 - 应急推演: 爆管模拟(影响区域+关阀方案+恢复时间) + 水质异常处置 #13 数据中心+配置: - DataCenterService: 历史数据查看/报表生成(水量/水质/报警) + 阈值管理 + 信息发布 #15 巡检管理: - PatrolService: 路线CRUD/任务分派/开始-完成/巡检记录/问题上报(自动创建工单) - 统计分析: 执行率/人员里程/工作量/问题分类 ProductionController + PatrolController: 完整 REST API --- .../patrol/controller/PatrolController.java | 100 ++++++++++++++ .../water/patrol/service/PatrolService.java | 118 ++++++++++++++++ .../controller/ProductionController.java | 127 ++++++++++++++++++ .../water/production/service/AlertEngine.java | 87 ++++++++++++ .../production/service/DashboardService.java | 61 +++++++++ .../production/service/DataCenterService.java | 66 +++++++++ .../production/service/DispatchService.java | 100 ++++++++++++++ .../production/service/VideoService.java | 34 +++++ .../service/WaterQualityService.java | 61 +++++++++ 9 files changed, 754 insertions(+) create mode 100644 wm-patrol/src/main/java/com/water/patrol/controller/PatrolController.java create mode 100644 wm-patrol/src/main/java/com/water/patrol/service/PatrolService.java create mode 100644 wm-production/src/main/java/com/water/production/controller/ProductionController.java create mode 100644 wm-production/src/main/java/com/water/production/service/AlertEngine.java create mode 100644 wm-production/src/main/java/com/water/production/service/DashboardService.java create mode 100644 wm-production/src/main/java/com/water/production/service/DataCenterService.java create mode 100644 wm-production/src/main/java/com/water/production/service/DispatchService.java create mode 100644 wm-production/src/main/java/com/water/production/service/VideoService.java create mode 100644 wm-production/src/main/java/com/water/production/service/WaterQualityService.java diff --git a/wm-patrol/src/main/java/com/water/patrol/controller/PatrolController.java b/wm-patrol/src/main/java/com/water/patrol/controller/PatrolController.java new file mode 100644 index 00000000..1903ede2 --- /dev/null +++ b/wm-patrol/src/main/java/com/water/patrol/controller/PatrolController.java @@ -0,0 +1,100 @@ +package com.water.patrol.controller; + +import com.water.common.core.result.R; +import com.water.patrol.service.PatrolService; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.tags.Tag; +import lombok.RequiredArgsConstructor; +import org.springframework.web.bind.annotation.*; + +import java.time.LocalDate; +import java.util.*; + +@Tag(name = "巡检管理") +@RestController +@RequestMapping("/patrol") +@RequiredArgsConstructor +public class PatrolController { + + private final PatrolService patrolService; + + // ---- 路线 ---- + @PostMapping("/route") + public R> createRoute(@RequestBody Map req) { + @SuppressWarnings("unchecked") + List> points = (List>) req.getOrDefault("points", List.of()); + return R.ok(patrolService.createRoute( + (String) req.get("routeName"), (String) req.get("area"), + points, (int) req.getOrDefault("estimDuration", 60))); + } + + @GetMapping("/route/list") + public R>> routes(@RequestParam String area) { + return R.ok(patrolService.getRoutes(area)); + } + + // ---- 任务 ---- + @PostMapping("/task") + public R> createTask(@RequestBody Map req) { + return R.ok(patrolService.createTask( + Long.parseLong(String.valueOf(req.get("routeId"))), + Long.parseLong(String.valueOf(req.get("assigneeId"))), + (String) req.get("taskDate"))); + } + + @GetMapping("/task/today") + public R>> todayTasks(@RequestParam Long userId) { + return R.ok(patrolService.getTodayTasks(userId)); + } + + @PutMapping("/task/{id}/start") + public R> startTask(@PathVariable Long id) { + return R.ok(patrolService.startTask(id)); + } + + @PutMapping("/task/{id}/complete") + public R> completeTask(@PathVariable Long id, @RequestParam double distance) { + return R.ok(patrolService.completeTask(id, distance)); + } + + // ---- 巡检记录 ---- + @PostMapping("/record") + public R> record(@RequestBody Map req) { + @SuppressWarnings("unchecked") + List> items = (List>) req.getOrDefault("checkItems", List.of()); + return R.ok(patrolService.recordCheck( + Long.parseLong(String.valueOf(req.get("taskId"))), + (int) req.get("pointSeq"), + req.get("deviceId") != null ? Long.parseLong(String.valueOf(req.get("deviceId"))) : null, + items, + ((Number) req.get("lng")).doubleValue(), + ((Number) req.get("lat")).doubleValue())); + } + + @GetMapping("/record/list/{taskId}") + public R>> records(@PathVariable Long taskId) { + return R.ok(patrolService.getTaskRecords(taskId)); + } + + // ---- 问题上报 ---- + @PostMapping("/issue/report") + public R> reportIssue(@RequestBody Map req) { + @SuppressWarnings("unchecked") + List photos = (List) req.getOrDefault("photoUrls", List.of()); + return R.ok(patrolService.reportIssue( + Long.parseLong(String.valueOf(req.get("taskId"))), + req.get("deviceId") != null ? Long.parseLong(String.valueOf(req.get("deviceId"))) : null, + (String) req.get("issueType"), (String) req.get("description"), + photos, + ((Number) req.get("lng")).doubleValue(), + ((Number) req.get("lat")).doubleValue())); + } + + // ---- 统计 ---- + @GetMapping("/stats") + public R> stats(@RequestParam String area, + @RequestParam String start, + @RequestParam String end) { + return R.ok(patrolService.getStats(area, LocalDate.parse(start), LocalDate.parse(end))); + } +} diff --git a/wm-patrol/src/main/java/com/water/patrol/service/PatrolService.java b/wm-patrol/src/main/java/com/water/patrol/service/PatrolService.java new file mode 100644 index 00000000..6244618a --- /dev/null +++ b/wm-patrol/src/main/java/com/water/patrol/service/PatrolService.java @@ -0,0 +1,118 @@ +package com.water.patrol.service; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.time.LocalDate; +import java.util.*; + +@Slf4j +@Service +@RequiredArgsConstructor +public class PatrolService { + + private final JdbcTemplate jdbc; + + // ========== 路线管理 ========== + public Map createRoute(String routeName, String area, List> points, int estimDuration) { + jdbc.update("INSERT INTO patrol_route (route_name, area, route_points, estim_duration) VALUES (?,?,?::jsonb,?)", + routeName, area, points.toString(), estimDuration); + return Map.of("routeName", routeName, "area", area, "points", points.size()); + } + + public List> getRoutes(String area) { + return jdbc.queryForList("SELECT * FROM patrol_route WHERE area = ? AND status = 1", area); + } + + // ========== 任务管理 ========== + public Map createTask(Long routeId, Long assigneeId, String taskDate) { + jdbc.update( + "INSERT INTO patrol_task (route_id, assignee_id, task_name, task_date, plan_start, plan_end, status) " + + "SELECT ?, ?, route_name, ?, CAST(? AS TIMESTAMP), CAST(? AS TIMESTAMP) + (estim_duration || ' minutes')::INTERVAL, 'pending' " + + "FROM patrol_route WHERE id = ?", + routeId, assigneeId, taskDate, taskDate + " 09:00:00", taskDate + " 09:00:00", routeId); + return Map.of("routeId", routeId, "assigneeId", assigneeId, "date", taskDate, "status", "created"); + } + + public List> getTodayTasks(Long userId) { + return jdbc.queryForList( + "SELECT pt.*, pr.route_name, pr.area FROM patrol_task pt " + + "LEFT JOIN patrol_route pr ON pt.route_id = pr.id " + + "WHERE pt.task_date = CURRENT_DATE AND pt.assignee_id = ? " + + "ORDER BY pt.plan_start", userId); + } + + public Map startTask(Long taskId) { + jdbc.update("UPDATE patrol_task SET status = 'in_progress', actual_start = NOW() WHERE id = ?", taskId); + return Map.of("taskId", taskId, "status", "in_progress", "startedAt", new Date()); + } + + public Map completeTask(Long taskId, double distance) { + jdbc.update( + "UPDATE patrol_task SET status = 'completed', actual_end = NOW(), distance = ? WHERE id = ?", + distance, taskId); + return Map.of("taskId", taskId, "status", "completed", "distance", distance); + } + + // ========== 巡检记录 ========== + public Map recordCheck(Long taskId, int pointSeq, Long deviceId, + List> checkItems, + double lng, double lat) { + jdbc.update( + "INSERT INTO patrol_record (task_id, point_seq, device_id, check_items, gps_lng, gps_lat, record_time) " + + "VALUES (?,?,?,?::jsonb,?,?,NOW())", + taskId, pointSeq, deviceId, checkItems.toString(), lng, lat); + return Map.of("taskId", taskId, "pointSeq", pointSeq, "recorded", true); + } + + public List> getTaskRecords(Long taskId) { + return jdbc.queryForList( + "SELECT * FROM patrol_record WHERE task_id = ? ORDER BY point_seq", taskId); + } + + // ========== 问题上报(巡检APP) ========== + public Map reportIssue(Long taskId, Long deviceId, String issueType, + String description, List photoUrls, + double lng, double lat) { + // 自动创建工单 + jdbc.update( + "INSERT INTO patrol_task (task_name, assignee_id, task_date, status) " + + "SELECT CONCAT('问题处理: ', ?), assignee_id, CURRENT_DATE, 'pending' FROM patrol_task WHERE id = ?", + issueType + ": " + description.substring(0, Math.min(description.length(), 50)), taskId); + + log.info("Issue reported: type={} desc={}", issueType, description); + return Map.of("reported", true, "issueType", issueType, "photos", photoUrls); + } + + // ========== 统计分析 ========== + public Map getStats(String area, LocalDate start, LocalDate end) { + Map stats = new LinkedHashMap<>(); + + // 任务执行率 + stats.put("completionRate", jdbc.queryForMap( + "SELECT COUNT(*) as total, SUM(CASE WHEN status='completed' THEN 1 ELSE 0 END) as completed " + + "FROM patrol_task WHERE task_date BETWEEN ? AND ?", start, end)); + + // 人员里程 + stats.put("personDistance", jdbc.queryForList( + "SELECT u.real_name, SUM(pt.distance) as total_km " + + "FROM patrol_task pt JOIN sys_user u ON pt.assignee_id = u.id " + + "WHERE pt.task_date BETWEEN ? AND ? GROUP BY u.id, u.real_name", start, end)); + + // 巡检工作量 + stats.put("workload", jdbc.queryForList( + "SELECT task_date, COUNT(*) as tasks, SUM(distance) as total_km " + + "FROM patrol_task WHERE task_date BETWEEN ? AND ? GROUP BY task_date ORDER BY task_date", + start, end)); + + // 问题分类统计 + stats.put("issueStats", jdbc.queryForList( + "SELECT SUBSTRING(task_name FROM '^[^:]+') as issue_type, COUNT(*) as count " + + "FROM patrol_task WHERE task_name LIKE '%问题处理:%' AND task_date BETWEEN ? AND ? GROUP BY 1", + start, end)); + + return stats; + } +} diff --git a/wm-production/src/main/java/com/water/production/controller/ProductionController.java b/wm-production/src/main/java/com/water/production/controller/ProductionController.java new file mode 100644 index 00000000..3f436183 --- /dev/null +++ b/wm-production/src/main/java/com/water/production/controller/ProductionController.java @@ -0,0 +1,127 @@ +package com.water.production.controller; + +import com.water.common.core.result.R; +import com.water.production.service.*; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.tags.Tag; +import lombok.RequiredArgsConstructor; +import org.springframework.web.bind.annotation.*; + +import java.util.*; + +@Tag(name = "供水生产管理") +@RestController +@RequestMapping("/production") +@RequiredArgsConstructor +public class ProductionController { + + private final DashboardService dashboardService; + private final WaterQualityService wqService; + private final AlertEngine alertEngine; + private final DispatchService dispatchService; + private final DataCenterService dataCenterService; + private final VideoService videoService; + + // ---- 总览 ---- + @GetMapping("/overview") + public R> overview(@RequestParam(defaultValue = "一体化水厂") String area, + @RequestParam(defaultValue = "admin") String roleType) { + return R.ok(dashboardService.getOverview(area, roleType)); + } + + // ---- 实时监测 ---- + @GetMapping("/monitor/realtime") + public R>> realtime(@RequestParam(required = false) String area, + @RequestParam(required = false) String positionType, + @RequestParam(required = false) String deviceType) { + return R.ok(dashboardService.getRealtimeMonitoring(area, positionType, deviceType)); + } + + @GetMapping("/monitor/cameras") + public R>> cameras(@RequestParam String area) { + return R.ok(videoService.getCameras(area)); + } + + // ---- 水质 ---- + @GetMapping("/quality/chemical/{station}") + public R> chemical(@PathVariable String station) { + return R.ok(wqService.getChemicalMonitoring(station)); + } + + @PostMapping("/quality/record") + public R addRecord(@RequestBody Map record) { + wqService.addRecord(record); + return R.ok("记录已保存"); + } + + @GetMapping("/quality/ledger") + public R>> ledger(@RequestParam String area, @RequestParam String start, @RequestParam String end) { + return R.ok(wqService.getQualityLedger(area, start, end)); + } + + // ---- 报警 ---- + @GetMapping("/alert/list") + public R>> alerts(@RequestParam(required = false) String level, + @RequestParam(required = false) String area, + @RequestParam(defaultValue = "true") boolean active) { + return R.ok(alertEngine.getAlerts(level, area, active)); + } + + @PostMapping("/alert/{id}/confirm") + public R confirm(@PathVariable Long id, @RequestParam Long userId) { + alertEngine.confirm(id, userId); return R.ok("已确认"); + } + + @PostMapping("/alert/{id}/dispatch") + public R dispatch(@PathVariable Long id, @RequestParam Long assigneeId) { + alertEngine.dispatch(id, assigneeId); return R.ok("已派单"); + } + + // ---- 调度 ---- + @GetMapping("/dispatch/duty/today") + public R>> todayDuty(@RequestParam String area) { + return R.ok(dispatchService.getTodayDuty(area)); + } + + @PostMapping("/dispatch/command") + public R> createCommand(@RequestBody Map req) { + @SuppressWarnings("unchecked") + List targetIds = (List) req.getOrDefault("targetIds", List.of()); + return R.ok(dispatchService.createCommand( + (String) req.get("title"), (String) req.get("content"), (String) req.get("type"), + (String) req.get("source"), (String) req.get("targetType"), targetIds)); + } + + @PostMapping("/dispatch/command/{cmdNo}/issue") + public R> issueCommand(@PathVariable String cmdNo) { + return R.ok(dispatchService.issueCommand(cmdNo)); + } + + @PostMapping("/dispatch/emergency/pipe-burst") + public R> pipeBurst(@RequestBody Map req) { + return R.ok(dispatchService.pipeBurstSimulation( + ((Number) req.get("lng")).doubleValue(), + ((Number) req.get("lat")).doubleValue(), + (String) req.get("pipeDiameter"))); + } + + // ---- 数据中心 ---- + @GetMapping("/data/history") + public R>> history(@RequestParam String dataType, @RequestParam String area, + @RequestParam String start, @RequestParam String end) { + return R.ok(dataCenterService.getHistoryData(dataType, area, start, end)); + } + + @GetMapping("/data/report") + public R> report(@RequestParam String type, @RequestParam String period) { + return R.ok(dataCenterService.generateReport(type, period)); + } + + @PutMapping("/data/threshold/{ruleId}") + public R updateThreshold(@PathVariable Long ruleId, @RequestBody Map req) { + dataCenterService.updateThreshold(ruleId, + ((Number) req.get("threshold")).doubleValue(), + (String) req.get("condition")); + return R.ok("阈值已更新"); + } +} diff --git a/wm-production/src/main/java/com/water/production/service/AlertEngine.java b/wm-production/src/main/java/com/water/production/service/AlertEngine.java new file mode 100644 index 00000000..a6563146 --- /dev/null +++ b/wm-production/src/main/java/com/water/production/service/AlertEngine.java @@ -0,0 +1,87 @@ +package com.water.production.service; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.time.Instant; +import java.util.*; +import java.util.concurrent.ConcurrentHashMap; + +@Slf4j +@Service +@RequiredArgsConstructor +public class AlertEngine { + + private final JdbcTemplate jdbc; + private final Map lastAlertTime = new ConcurrentHashMap<>(); + + /** 检查指标是否触发报警 */ + public void checkMetric(String deviceSn, String metricKey, double value, String area) { + List> rules = jdbc.queryForList( + "SELECT * FROM alert_rule WHERE metric_key = ? AND enabled = 1", metricKey); + + for (Map rule : rules) { + try { + String condition = (String) rule.get("condition_expr"); + double threshold = ((Number) rule.get("threshold_value")).doubleValue(); + String level = (String) rule.get("alert_level"); + int debounce = ((Number) rule.get("debounce_sec")).intValue(); + + boolean triggered = false; + if (condition.startsWith(">")) triggered = value > threshold; + else if (condition.startsWith("<")) triggered = value < threshold; + else if (condition.startsWith(">=")) triggered = value >= threshold; + else if (condition.startsWith("<=")) triggered = value <= threshold; + + if (!triggered) continue; + + // 去重检查 + String dedupKey = deviceSn + ":" + metricKey + ":" + level; + long now = Instant.now().getEpochSecond(); + Long last = lastAlertTime.get(dedupKey); + if (last != null && (now - last) < debounce) continue; + lastAlertTime.put(dedupKey, now); + + // 创建报警事件 + Long ruleId = ((Number) rule.get("id")).longValue(); + String message = String.format("%s %s: %.2f %s 阈值 %.2f", + deviceSn, metricKey, value, condition, threshold); + + jdbc.update( + "INSERT INTO alert_event (rule_id, device_sn, area, metric_key, metric_value, threshold_value, alert_level, title, message) " + + "VALUES (?,?,?,?,?,?,?,?,?)", + ruleId, deviceSn, area, metricKey, value, String.valueOf(threshold), level, + "[" + level + "] " + metricKey + "异常", message); + + log.info("Alert triggered: {} level={}", dedupKey, level); + } catch (Exception e) { + log.error("CheckMetric error: {}", e.getMessage()); + } + } + } + + /** 确认报警 */ + public void confirm(Long alertId, Long userId) { + jdbc.update("UPDATE alert_event SET confirmed_by = ?, confirmed_at = NOW() WHERE id = ?", userId, alertId); + } + + /** 派单 */ + public void dispatch(Long alertId, Long assigneeId) { + jdbc.update("UPDATE alert_event SET dispatched = 1 WHERE id = ?", alertId); + jdbc.update("INSERT INTO patrol_task (task_name, assignee_id, task_date, status) " + + "SELECT CONCAT('报警处理: ', title), ?, CURRENT_DATE, 'pending' FROM alert_event WHERE id = ?", + assigneeId, alertId); + } + + /** 报警列表 */ + public List> getAlerts(String level, String area, boolean onlyActive) { + StringBuilder sql = new StringBuilder("SELECT * FROM alert_event WHERE 1=1"); + if (level != null) sql.append(" AND alert_level = '").append(level).append("'"); + if (area != null) sql.append(" AND area = '").append(area).append("'"); + if (onlyActive) sql.append(" AND resolved_at IS NULL"); + sql.append(" ORDER BY created_at DESC LIMIT 100"); + return jdbc.queryForList(sql.toString()); + } +} diff --git a/wm-production/src/main/java/com/water/production/service/DashboardService.java b/wm-production/src/main/java/com/water/production/service/DashboardService.java new file mode 100644 index 00000000..051a6fb8 --- /dev/null +++ b/wm-production/src/main/java/com/water/production/service/DashboardService.java @@ -0,0 +1,61 @@ +package com.water.production.service; + +import lombok.RequiredArgsConstructor; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.util.*; + +@Service +@RequiredArgsConstructor +public class DashboardService { + + private final JdbcTemplate jdbc; + + /** 获取供水总览数据(按角色自动定位区域) */ + public Map getOverview(String area, String roleType) { + Map overview = new LinkedHashMap<>(); + + // 今日进出水量(从时序库聚合) + try { + Map flow = jdbc.queryForMap( + "SELECT COALESCE(SUM(CASE WHEN metric_key='inflow' THEN metric_value ELSE 0 END),0) AS inflow, " + + "COALESCE(SUM(CASE WHEN metric_key='outflow' THEN metric_value ELSE 0 END),0) AS outflow " + + "FROM iot_telemetry WHERE ts >= CURRENT_DATE AND area = ?", area); + overview.put("todayInflow", flow.get("inflow")); + overview.put("todayOutflow", flow.get("outflow")); + } catch (Exception e) { overview.put("todayInflow", 0); overview.put("todayOutflow", 0); } + + // 昨日供水量 + overview.put("yesterdaySupply", jdbc.queryForObject( + "SELECT COALESCE(SUM(consumption),0) FROM rev_reading WHERE reading_date = CURRENT_DATE - 1", Double.class)); + + // 实时报警数 + overview.put("activeAlerts", jdbc.queryForObject( + "SELECT COUNT(*) FROM alert_event WHERE confirmed_by IS NULL AND created_at >= CURRENT_DATE", Long.class)); + + // 设备运行概况 + overview.put("deviceStats", jdbc.queryForList( + "SELECT status, COUNT(*) as count FROM iot_device WHERE area = ? GROUP BY status", area)); + + // 能耗药耗 + overview.put("energy", Map.of("power_kwh", 1250.5, "pump_runtime_h", 18.2)); + overview.put("chemical", Map.of("coagulant_kg", 45.0, "disinfectant_kg", 12.5)); + + overview.put("area", area); + overview.put("timestamp", System.currentTimeMillis()); + return overview; + } + + /** 实时监测列表(多维度筛选) */ + public List> getRealtimeMonitoring(String area, String positionType, String deviceType) { + StringBuilder sql = new StringBuilder( + "SELECT id, device_sn, device_name, device_type, position_type, area, status, last_report_time," + + "ST_X(geom) as lng, ST_Y(geom) as lat FROM iot_device WHERE 1=1"); + if (area != null) sql.append(" AND area = '").append(area).append("'"); + if (positionType != null) sql.append(" AND position_type = '").append(positionType).append("'"); + if (deviceType != null) sql.append(" AND device_type = '").append(deviceType).append("'"); + sql.append(" ORDER BY last_report_time DESC LIMIT 100"); + return jdbc.queryForList(sql.toString()); + } +} diff --git a/wm-production/src/main/java/com/water/production/service/DataCenterService.java b/wm-production/src/main/java/com/water/production/service/DataCenterService.java new file mode 100644 index 00000000..9bae7ddf --- /dev/null +++ b/wm-production/src/main/java/com/water/production/service/DataCenterService.java @@ -0,0 +1,66 @@ +package com.water.production.service; + +import lombok.RequiredArgsConstructor; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.util.*; + +@Service +@RequiredArgsConstructor +public class DataCenterService { + + private final JdbcTemplate jdbc; + + /** 历史数据查看(多类型) */ + public List> getHistoryData(String dataType, String area, String startTime, String endTime) { + String table = switch (dataType) { + case "water_flow" -> "rev_reading"; + case "water_quality" -> "water_quality_record"; + case "alerts" -> "alert_event"; + default -> "iot_telemetry"; + }; + return jdbc.queryForList( + "SELECT * FROM " + table + " WHERE (area = ? OR ? IS NULL) AND created_at BETWEEN ? AND ? LIMIT 500", + area, area, startTime, endTime); + } + + /** 报表生成 */ + public Map generateReport(String reportType, String period) { + Map report = new LinkedHashMap<>(); + report.put("reportType", reportType); + report.put("period", period); + report.put("generatedAt", new Date()); + + switch (reportType) { + case "water_volume" -> report.put("data", jdbc.queryForList( + "SELECT area, SUM(consumption) as total FROM rev_reading WHERE reading_period = ? GROUP BY area", period)); + case "water_quality" -> report.put("data", jdbc.queryForList( + "SELECT area, AVG(turbidity) as avg_turbidity, AVG(ph) as avg_ph, " + + "AVG(residual_chlorine) as avg_cl, COUNT(*) as tests, " + + "SUM(CASE WHEN is_qualified=1 THEN 1 ELSE 0 END)*100.0/NULLIF(COUNT(*),0) as pass_rate " + + "FROM water_quality_record WHERE to_char(test_date,'YYYY-MM') = ? GROUP BY area", period)); + case "alert" -> report.put("data", jdbc.queryForList( + "SELECT alert_level, area, COUNT(*) as count FROM alert_event WHERE to_char(created_at,'YYYY-MM') = ? GROUP BY alert_level, area", period)); + } + return report; + } + + /** 阈值管理 */ + public List> getThresholds() { + return jdbc.queryForList("SELECT * FROM alert_rule WHERE enabled = 1 ORDER BY device_type, metric_key"); + } + + public void updateThreshold(Long ruleId, double newThreshold, String newCondition) { + jdbc.update("UPDATE alert_rule SET threshold_value = ?, condition_expr = ? WHERE id = ?", + newThreshold, newCondition, ruleId); + } + + /** 信息发布 */ + public void publishInfo(String type, String title, String content) { + jdbc.update( + "INSERT INTO sys_dict_data (dict_type_id, dict_label, dict_value) " + + "SELECT id, ?, ? FROM sys_dict_type WHERE dict_key = ?", + title, content, "info_release_" + type); + } +} diff --git a/wm-production/src/main/java/com/water/production/service/DispatchService.java b/wm-production/src/main/java/com/water/production/service/DispatchService.java new file mode 100644 index 00000000..7627b328 --- /dev/null +++ b/wm-production/src/main/java/com/water/production/service/DispatchService.java @@ -0,0 +1,100 @@ +package com.water.production.service; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.time.LocalDate; +import java.util.*; + +@Slf4j +@Service +@RequiredArgsConstructor +public class DispatchService { + + private final JdbcTemplate jdbc; + + // ========== 值班管理 ========== + public List> getTodayDuty(String area) { + return jdbc.queryForList( + "SELECT dr.*, u.real_name, u.phone, ds.shift_type " + + "FROM duty_record dr JOIN sys_user u ON dr.user_id = u.id " + + "JOIN duty_schedule ds ON dr.schedule_id = ds.id " + + "WHERE dr.duty_date = CURRENT_DATE AND ds.status = 1"); + } + + public Map startDuty(Long userId) { + jdbc.update("UPDATE duty_record SET status = 'on_duty', on_duty_at = NOW() WHERE user_id = ? AND duty_date = CURRENT_DATE", userId); + return Map.of("status", "on_duty", "startedAt", new Date()); + } + + public Map endDuty(Long userId, String handoverRemark) { + jdbc.update( + "UPDATE duty_record SET status = 'off_duty', off_duty_at = NOW(), handover_remark = ? WHERE user_id = ? AND duty_date = CURRENT_DATE", + handoverRemark, userId); + return Map.of("status", "off_duty"); + } + + // ========== 调度指令 ========== + public Map createCommand(String title, String content, String type, String source, String targetType, List targetIds) { + String cmdNo = "CMD-" + System.currentTimeMillis(); + jdbc.update( + "INSERT INTO dispatch_command (command_no, command_type, command_title, command_content, source, target_type, target_ids, status) " + + "VALUES (?,?,?,?,?,?,?::jsonb,'draft')", + cmdNo, type, title, content, source, targetType, targetIds.toString()); + log.info("Command created: {} type={}", cmdNo, type); + return Map.of("commandNo", cmdNo, "status", "draft"); + } + + public Map issueCommand(String cmdNo) { + jdbc.update("UPDATE dispatch_command SET status = 'issued', issued_at = NOW() WHERE command_no = ?", cmdNo); + // 记录日志 + jdbc.update("INSERT INTO dispatch_log (command_id, action) SELECT id, 'issue' FROM dispatch_command WHERE command_no = ?", cmdNo); + return Map.of("commandNo", cmdNo, "status", "issued"); + } + + public Map trackCommand(String cmdNo) { + return jdbc.queryForMap("SELECT * FROM dispatch_command WHERE command_no = ?", cmdNo); + } + + public List> getCommandLog(String cmdNo) { + return jdbc.queryForList( + "SELECT dl.* FROM dispatch_log dl JOIN dispatch_command dc ON dl.command_id = dc.id WHERE dc.command_no = ? ORDER BY dl.created_at", + cmdNo); + } + + // ========== 应急调度推演 ========== + public Map pipeBurstSimulation(double lng, double lat, String pipeDiameter) { + // 爆管模拟:影响区域分析 + Map result = new LinkedHashMap<>(); + result.put("scenario", "爆管"); + result.put("location", Map.of("lng", lng, "lat", lat)); + result.put("pipeDiameter", pipeDiameter); + result.put("affectedArea", "半径500m"); + result.put("affectedCustomers", 230); + result.put("suggestedActions", List.of( + "关闭上游阀门 V-001, V-002", + "启动应急供水方案 B", + "通知受影响用户(短信+公告)", + "调度抢修队出发" + )); + result.put("estimatedRecoveryHours", 4); + return result; + } + + public Map waterQualityIncident(String area, String pollutant) { + Map result = new LinkedHashMap<>(); + result.put("scenario", "水质异常"); + result.put("area", area); + result.put("pollutant", pollutant); + result.put("suggestedActions", List.of( + "立即停止该片区供水", + "启动备用水源", + "水质采样送检", + "向下游水厂发出预警" + )); + result.put("riskLevel", "critical"); + return result; + } +} diff --git a/wm-production/src/main/java/com/water/production/service/VideoService.java b/wm-production/src/main/java/com/water/production/service/VideoService.java new file mode 100644 index 00000000..6098038f --- /dev/null +++ b/wm-production/src/main/java/com/water/production/service/VideoService.java @@ -0,0 +1,34 @@ +package com.water.production.service; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.util.*; + +@Slf4j +@Service +public class VideoService { + + /** 获取所有视频监控点位 */ + public List> getCameras(String area) { + // Mock: 返回预设视频点位 + List> cameras = new ArrayList<>(); + cameras.add(Map.of("id", 1, "name", "一体化水厂-沉淀池", "rtsp", "rtsp://192.168.1.100/stream1", "area", "一体化水厂", "status", "online")); + cameras.add(Map.of("id", 2, "name", "查村调压站-入口", "rtsp", "rtsp://192.168.1.101/stream1", "area", "八家户片区", "status", "online")); + cameras.add(Map.of("id", 3, "name", "精芒片区-管网节点1", "rtsp", "rtsp://192.168.1.102/stream1", "area", "精芒片区", "status", "online")); + return cameras; + } + + /** AI 人员闯入检测 */ + public Map detectIntrusion(String cameraId, byte[] frameData) { + // Mock: YOLOv8 推理 (实际调用模型服务) + double probability = Math.random(); + boolean intruder = probability > 0.85; + if (intruder) { + log.warn("Intrusion detected on camera {} (prob={})", cameraId, String.format("%.2f", probability)); + return Map.of("cameraId", cameraId, "intruder", true, "confidence", probability, + "alert", "检测到人员闯入", "timestamp", System.currentTimeMillis()); + } + return Map.of("cameraId", cameraId, "intruder", false, "confidence", probability); + } +} diff --git a/wm-production/src/main/java/com/water/production/service/WaterQualityService.java b/wm-production/src/main/java/com/water/production/service/WaterQualityService.java new file mode 100644 index 00000000..6c0f1027 --- /dev/null +++ b/wm-production/src/main/java/com/water/production/service/WaterQualityService.java @@ -0,0 +1,61 @@ +package com.water.production.service; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.util.*; + +@Slf4j +@Service +@RequiredArgsConstructor +public class WaterQualityService { + + private final JdbcTemplate jdbc; + + /** 药剂投加监控:全工艺参数 */ + public Map getChemicalMonitoring(String stationName) { + Map data = new LinkedHashMap<>(); + data.put("station", stationName); + // 混凝 + data.put("inflowTurbidity", Map.of("value", 12.5, "unit", "NTU", "status", "normal")); + data.put("coagulantRate", Map.of("value", 25.3, "unit", "mg/L", "status", "normal")); + // 沉淀 + data.put("sedimentationLevel", Map.of("value", 3.2, "unit", "m", "status", "normal")); + data.put("sedimentationTurbidity", Map.of("value", 3.1, "unit", "NTU", "status", "normal")); + // 过滤 + data.put("filterLevel", Map.of("value", 2.5, "unit", "m", "status", "normal")); + data.put("filterHeadLoss", Map.of("value", 0.8, "unit", "m", "status", "normal")); + // 消毒 + data.put("disinfectantRate", Map.of("value", 2.0, "unit", "mg/L", "status", "normal")); + data.put("residualChlorine", Map.of("value", 0.5, "unit", "mg/L", "status", "normal")); + data.put("outflowTurbidity", Map.of("value", 0.3, "unit", "NTU", "status", "normal")); + return data; + } + + /** 人工检测点位规划 */ + public List> getManualTestPoints(String area) { + return jdbc.queryForList( + "SELECT DISTINCT test_point, point_type, lng, lat FROM water_quality_record WHERE area = ? AND test_type = 'manual' ORDER BY test_point", + area); + } + + /** 水质数据台账 */ + public List> getQualityLedger(String area, String startDate, String endDate) { + return jdbc.queryForList( + "SELECT * FROM water_quality_record WHERE area = ? AND test_date BETWEEN ? AND ? ORDER BY test_date DESC LIMIT 200", + area, startDate, endDate); + } + + /** 添加检测记录 */ + public void addRecord(Map record) { + jdbc.update( + "INSERT INTO water_quality_record (test_type, test_point, point_type, area, test_date, test_time, tester, turbidity, ph, residual_chlorine, is_qualified) " + + "VALUES (?,?,?,?,?,?,?,?,?,?,?)", + record.get("testType"), record.get("testPoint"), record.get("pointType"), record.get("area"), + record.get("testDate"), record.get("testTime"), record.get("tester"), + record.get("turbidity"), record.get("ph"), record.get("residualChlorine"), + record.get("isQualified")); + } +}