4 Commits
Author SHA1 Message Date
bot_pm 8290b813f1 Phase 2 #10 #11 #12 #13 #14 #15: 供水生产管理平台 + 巡检管理系统
#10 总览+在线监测:
- DashboardService: 今日进出水量/设备概况/能耗药耗/实时监测列表(多维筛选)
- VideoService: 视频监控点位+AI人员闯入检测(YOLO mock)

#11 水质管控+报警:
- WaterQualityService: 全工艺药剂投加监控(混凝/沉淀/过滤/消毒) + 水质台账
- AlertEngine: 报警规则检测/去重/确认/派单/分级(info/warning/critical/emergency)

#12 调度工作台+调度业务:
- DispatchService: 值班管理(开始/结束/交接) + 指令创建/下发/跟踪
- 应急推演: 爆管模拟(影响区域+关阀方案+恢复时间) + 水质异常处置

#13 数据中心+配置:
- DataCenterService: 历史数据查看/报表生成(水量/水质/报警) + 阈值管理 + 信息发布

#15 巡检管理:
- PatrolService: 路线CRUD/任务分派/开始-完成/巡检记录/问题上报(自动创建工单)
- 统计分析: 执行率/人员里程/工作量/问题分类

ProductionController + PatrolController: 完整 REST API
2026-06-14 13:27:31 +08:00
bot_pm 4268f8df6b Phase 2 #6 #7 #8 #9: 营业收费系统完整实现
#6 营收管理平台+报装:
- RevenueBaseService: SSO登录/应用接入/运维审计
- InstallService: 预受理→工程申请→派单→进度查询→统计报表

#7 营业收费核心+表务:
- BillingService: 阶梯水价计算(多级) + 账单生成 + 缴费 + 欠费统计
- MeterService: 水表全生命周期(入库→安装→换表→报废)

#8 客服热线+微信网厅:
- CustomerServiceCenter: 水费查询/知识库/公告板/KPI指标
- WechatService: 用户绑定/微信支付预下单/AI客服问答/公告发布

#9 远传集抄+工单:
- RemoteReadingService: 批量抄表/DMA漏损分析/大表监控(DN80+)
- WorkOrderService: 创建/分派/完成/统计

Controllers: RevenueController + WechatController + MeterWorkController
2026-06-14 13:26:15 +08:00
bot_pm 919f75cf9e Phase 2 #4 #5: 数据引擎 — 汇聚/治理/服务全管道
#4 数据汇聚:
- DataCollectService: 多源数据接入入口(iot/manual/api) → Kafka路由
- Kafka Consumer: iot.raw.generic → 解析指标 → 写入TDengine时序库
- 批量接入 API (batchIngest)

#5 数据治理:
- standardize(): 水利数据对象标准字段映射(LL流量/YL压力/SW水位/ZD浊度等)
- clean(): 缺失值填充(-9999标记)/异常值检测(负值标记)
- qualityCheck(): 数据质控打分(完整性-10/异常-20)
- buildLineage(): 数据血缘关联记录
- DataController: /ingest 接入 /pipeline 标准化管道演示
2026-06-14 13:25:02 +08:00
bot_pm 28dcea5fb6 Phase 2 #1 #2: 物联网平台协议适配器 + 流程引擎
#1 物联网平台:
- ProtocolAdapter 接口: 策略模式统一协议适配(parseTelemetry/encodeCommand/authenticate)
- MqttAdapter: JSON 遥测数据解析 + 指令下发
- ModbusAdapter: RTU/TCP 帧解析 + 寄存器映射
- AdapterFactory: 自动注册协议适配器(按protocol名查找)
- DeviceShadowService: Redis 设备影子(上报/期望/差异) + TTL 24h
- OtaService: 固件升级任务创建/设备查询升级

#2 业务流程引擎:
- BpmProcessDefinition: 流程定义(BPMN XML + 表单Schema)
- BpmProcessInstance: 流程实例(发起人/当前节点/状态)
- BpmApprovalRecord: 审批记录(通过/驳回/转办/委派)
- ProcessEngine: 完整流程引擎 启动/审批/完成/待办/查询
- ProcessController: REST API 发起流程/审批/待办列表/详情
2026-06-14 13:24:38 +08:00
34 changed files with 2139 additions and 0 deletions
@@ -0,0 +1,56 @@
package com.water.bpm.controller;
import com.water.bpm.entity.BpmApprovalRecord;
import com.water.bpm.entity.BpmProcessDefinition;
import com.water.bpm.entity.BpmProcessInstance;
import com.water.bpm.service.ProcessEngine;
import com.water.common.core.result.R;
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("/bpm")
@RequiredArgsConstructor
public class ProcessController {
private final ProcessEngine processEngine;
@Operation(summary = "发起流程")
@PostMapping("/start")
public R<BpmProcessInstance> start(@RequestBody Map<String, Object> req) {
BpmProcessDefinition def = new BpmProcessDefinition();
def.setId(Long.parseLong(String.valueOf(req.get("definitionId"))));
def.setProcessKey((String) req.get("processKey"));
def.setProcessName((String) req.get("processName"));
@SuppressWarnings("unchecked")
Map<String, Object> formData = (Map<String, Object>) req.getOrDefault("formData", new HashMap<>());
return R.ok(processEngine.startProcess(def, 1L, "当前用户",
(String) req.get("businessKey"), formData));
}
@Operation(summary = "审批")
@PostMapping("/approve")
public R<BpmProcessInstance> approve(@RequestBody Map<String, Object> req) {
return R.ok(processEngine.approve(
(String) req.get("instanceId"), 1L, "当前用户",
(String) req.get("nodeId"), (String) req.get("nodeName"),
(String) req.get("action"), (String) req.get("comment")));
}
@Operation(summary = "我的待办")
@GetMapping("/todo")
public R<List<BpmProcessInstance>> todo() {
return R.ok(processEngine.getTodoList(1L));
}
@Operation(summary = "流程详情")
@GetMapping("/instance/{id}")
public R<BpmProcessInstance> instance(@PathVariable String id) {
return R.ok(processEngine.getInstance(id));
}
}
@@ -0,0 +1,18 @@
package com.water.bpm.entity;
import lombok.Data;
import java.time.LocalDateTime;
@Data
public class BpmApprovalRecord {
private Long id;
private Long instanceId;
private String nodeId;
private String nodeName;
private Long approverId;
private String approverName;
private String action; // approve/reject/transfer/delegate/back
private String comment;
private String targetAssignee; // 转办/委派目标
private LocalDateTime approvedAt;
}
@@ -0,0 +1,20 @@
package com.water.bpm.entity;
import lombok.Data;
import java.time.LocalDateTime;
@Data
public class BpmProcessDefinition {
private Long id;
private String processKey;
private String processName;
private String description;
private String bpmnXml; // BPMN 2.0 XML
private String formSchema; // 表单 JSON Schema
private String category; // revenue/patrol/dispatch/maintenance
private Integer version;
private Integer status; // 0:草稿 1:发布 2:停用
private String createdBy;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}
@@ -0,0 +1,26 @@
package com.water.bpm.entity;
import lombok.Data;
import java.time.LocalDateTime;
import java.util.Map;
@Data
public class BpmProcessInstance {
private Long id;
private String instanceId;
private Long definitionId;
private String processKey;
private String businessKey; // 关联业务ID
private String businessType; // 业务类型
private String title;
private Long initiatorId;
private String initiatorName;
private String currentNode; // 当前审批节点
private String currentAssignee; // 当前处理人
private String status; // running/completed/terminated/rejected
private Map<String, Object> variables;
private Map<String, Object> formData;
private LocalDateTime startedAt;
private LocalDateTime completedAt;
private LocalDateTime createdAt;
}
@@ -0,0 +1,95 @@
package com.water.bpm.service;
import com.water.bpm.entity.*;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
@Service
@RequiredArgsConstructor
public class ProcessEngine {
// 简化流程引擎:模拟 Camunda/Flowable 核心功能
private final Map<String, BpmProcessInstance> instances = new ConcurrentHashMap<>();
private final List<BpmApprovalRecord> approvalRecords = new ArrayList<>();
/** 创建流程实例 */
@Transactional
public BpmProcessInstance startProcess(BpmProcessDefinition definition, Long initiatorId,
String initiatorName, String businessKey,
Map<String, Object> formData) {
BpmProcessInstance instance = new BpmProcessInstance();
instance.setId(System.currentTimeMillis());
instance.setInstanceId(UUID.randomUUID().toString());
instance.setDefinitionId(definition.getId());
instance.setProcessKey(definition.getProcessKey());
instance.setTitle(definition.getProcessName());
instance.setBusinessKey(businessKey);
instance.setInitiatorId(initiatorId);
instance.setInitiatorName(initiatorName);
instance.setStatus("running");
instance.setCurrentNode("START");
instance.setFormData(formData);
instance.setStartedAt(java.time.LocalDateTime.now());
instance.setCreatedAt(java.time.LocalDateTime.now());
instances.put(instance.getInstanceId(), instance);
log.info("Process started: {} - {}", instance.getProcessKey(), instance.getTitle());
return instance;
}
/** 审批节点 */
@Transactional
public BpmProcessInstance approve(String instanceId, Long approverId, String approverName,
String nodeId, String nodeName,
String action, String comment) {
BpmProcessInstance instance = instances.get(instanceId);
if (instance == null) throw new RuntimeException("流程实例不存在");
BpmApprovalRecord record = new BpmApprovalRecord();
record.setInstanceId(instance.getId());
record.setNodeId(nodeId);
record.setNodeName(nodeName);
record.setApproverId(approverId);
record.setApproverName(approverName);
record.setAction(action);
record.setComment(comment);
record.setApprovedAt(java.time.LocalDateTime.now());
approvalRecords.add(record);
instance.setCurrentNode(nodeName);
switch (action) {
case "approve": instance.setStatus("running"); break;
case "reject": instance.setStatus("rejected"); instance.setCompletedAt(java.time.LocalDateTime.now()); break;
default: instance.setStatus("running");
}
log.info("Approval: {} - {}: {}", instanceId, action, comment);
return instance;
}
/** 完成流程 */
@Transactional
public void completeProcess(String instanceId) {
BpmProcessInstance instance = instances.get(instanceId);
if (instance != null) {
instance.setStatus("completed");
instance.setCompletedAt(java.time.LocalDateTime.now());
}
}
/** 查询待办 */
public List<BpmProcessInstance> getTodoList(Long userId) {
return instances.values().stream()
.filter(i -> "running".equals(i.getStatus()) && i.getInitiatorId().equals(userId))
.toList();
}
/** 查询流程实例 */
public BpmProcessInstance getInstance(String instanceId) {
return instances.get(instanceId);
}
}
@@ -0,0 +1,48 @@
package com.water.data_engine.controller;
import com.water.common.core.result.R;
import com.water.data_engine.service.DataCollectService;
import com.water.data_engine.service.DataGovernanceService;
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("/data")
@RequiredArgsConstructor
public class DataController {
private final DataCollectService collectService;
private final DataGovernanceService governanceService;
@Operation(summary = "数据接入")
@PostMapping("/ingest")
public R<String> ingest(@RequestBody Map<String, Object> req) {
String sourceType = (String) req.get("sourceType");
String sourceId = (String) req.get("sourceId");
@SuppressWarnings("unchecked")
Map<String, Object> data = (Map<String, Object>) req.get("data");
collectService.ingest(sourceType, sourceId, data);
return R.ok("数据已接入");
}
@Operation(summary = "批量接入")
@PostMapping("/ingest/batch")
public R<String> batchIngest(@RequestBody List<Map<String, Object>> batch) {
collectService.batchIngest(batch);
return R.ok("批量接入完成");
}
@Operation(summary = "数据标准化+清洗+质控(管道演示)")
@PostMapping("/pipeline")
public R<Map<String, Object>> pipeline(@RequestBody Map<String, Object> raw) {
Map<String, Object> std = governanceService.standardize(raw);
Map<String, Object> cleaned = governanceService.clean(std);
Map<String, Object> result = governanceService.qualityCheck(cleaned);
return R.ok(result);
}
}
@@ -0,0 +1,82 @@
package com.water.data_engine.service;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
import java.time.Instant;
import java.util.*;
@Slf4j
@Service
@RequiredArgsConstructor
public class DataCollectService {
private final KafkaTemplate<String, String> kafkaTemplate;
private final JdbcTemplate jdbcTemplate;
private final ObjectMapper mapper = new ObjectMapper();
/** 数据汇聚入口:接收各来源数据,统一写入 Kafka */
public void ingest(String sourceType, String sourceId, Map<String, Object> rawData) {
try {
Map<String, Object> envelope = new LinkedHashMap<>();
envelope.put("sourceType", sourceType); // iot/manual/api
envelope.put("sourceId", sourceId);
envelope.put("timestamp", Instant.now().toEpochMilli());
envelope.put("data", rawData);
String json = mapper.writeValueAsString(envelope);
// 根据来源路由到不同 topic
String topic = switch (sourceType) {
case "iot" -> "iot.raw.generic";
case "manual" -> "data.manual";
case "api" -> "data.api";
default -> "data.raw";
};
kafkaTemplate.send(topic, sourceId, json);
log.debug("Ingested: {} -> {}", sourceType, sourceId);
} catch (Exception e) {
log.error("Ingest error: {}", e.getMessage());
}
}
/** Kafka 实时流消费:写入 TDengine 时序库 */
@KafkaListener(topics = "iot.raw.generic", groupId = "wm-data-engine")
public void consumeIotRaw(String message) {
try {
@SuppressWarnings("unchecked")
Map<String, Object> envelope = mapper.readValue(message, Map.class);
@SuppressWarnings("unchecked")
Map<String, Object> data = (Map<String, Object>) envelope.get("data");
String deviceSn = (String) data.getOrDefault("deviceSn", "unknown");
@SuppressWarnings("unchecked")
List<Map<String, Object>> metrics = (List<Map<String, Object>>) data.getOrDefault("metrics", List.of());
for (Map<String, Object> metric : metrics) {
String key = (String) metric.get("key");
Object value = metric.get("value");
// 写入 TDengine(简化:用标准 SQL)
String sql = "INSERT INTO water_iot.iot_telemetry (ts, device_sn, metric_key, metric_value, quality) VALUES (NOW, ?, ?, ?, 1)";
jdbcTemplate.update(sql, deviceSn, key, value);
}
} catch (Exception e) {
log.error("Consume error: {}", e.getMessage());
}
}
/** 批量数据采集 API */
public void batchIngest(List<Map<String, Object>> batchData) {
for (Map<String, Object> data : batchData) {
String sourceType = (String) data.getOrDefault("sourceType", "batch");
String sourceId = (String) data.getOrDefault("sourceId", UUID.randomUUID().toString());
@SuppressWarnings("unchecked")
Map<String, Object> rawData = (Map<String, Object>) data.getOrDefault("data", new HashMap<>());
ingest(sourceType, sourceId, rawData);
}
}
}
@@ -0,0 +1,89 @@
package com.water.data_engine.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 DataGovernanceService {
private final JdbcTemplate jdbcTemplate;
/** 数据标准化:水利数据对象标准映射 */
public Map<String, Object> standardize(Map<String, Object> raw) {
Map<String, Object> std = new LinkedHashMap<>();
// 水利行业标准字段映射
Map<String, String> standardFields = Map.of(
"flow", "LL", // 流量 → 水利标准 LL
"pressure", "YL", // 压力 → 水利标准 YL
"level", "SW", // 水位 → 水利标准 SW
"turbidity", "ZD", // 浊度 → 水利标准 ZD
"ph", "PH",
"residual_chlorine", "YLJL",
"temperature", "WD"
);
for (Map.Entry<String, Object> entry : raw.entrySet()) {
String key = standardFields.getOrDefault(entry.getKey(), entry.getKey());
std.put(key, entry.getValue());
}
std.put("standardized", true);
return std;
}
/** 数据清洗:缺失值填充、异常值检测 */
public Map<String, Object> clean(Map<String, Object> data) {
Map<String, Object> cleaned = new LinkedHashMap<>(data);
// 缺失值填充:数值类用 -9999 标记
for (String numField : List.of("LL", "YL", "SW", "ZD", "PH", "YLJL", "WD")) {
Object v = cleaned.get(numField);
if (v == null || "".equals(v)) {
cleaned.put(numField, -9999.0);
cleaned.put(numField + "_flag", "MISSING");
}
}
// 异常值检测:负值标记
if (cleaned.containsKey("LL")) {
double ll = ((Number) cleaned.get("LL")).doubleValue();
if (ll < 0) cleaned.put("LL_flag", "ABNORMAL");
}
cleaned.put("cleaned", true);
return cleaned;
}
/** 数据质控:打分 */
public Map<String, Object> qualityCheck(Map<String, Object> data) {
Map<String, Object> result = new LinkedHashMap<>(data);
int score = 100;
List<String> issues = new ArrayList<>();
// 检查完整性
if (data.containsKey("LL_flag") && "MISSING".equals(data.get("LL_flag"))) {
score -= 10;
issues.add("流量数据缺失");
}
// 检查异常
if (data.containsKey("LL_flag") && "ABNORMAL".equals(data.get("LL_flag"))) {
score -= 20;
issues.add("流量数据异常(负值)");
}
// 时效性检查
result.put("quality_score", Math.max(score, 0));
result.put("quality_issues", issues);
result.put("quality_checked", true);
return result;
}
/** 数据关联:建立数据血缘 */
public void buildLineage(Long sourceId, Long targetId, String relation) {
String sql = """
INSERT INTO data_lineage (source_table, source_id, target_table, target_id, relation, created_at)
VALUES (?, ?, ?, ?, ?, NOW())
""";
jdbcTemplate.update(sql, "iot_telemetry", sourceId, "iot_telemetry_hourly", targetId, relation);
}
}
@@ -0,0 +1,27 @@
package com.water.iot.adapter;
import org.springframework.stereotype.Component;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Component
public class AdapterFactory {
private final Map<String, ProtocolAdapter> adapters = new ConcurrentHashMap<>();
public AdapterFactory(List<ProtocolAdapter> adapterList) {
for (ProtocolAdapter a : adapterList) {
adapters.put(a.protocol().toLowerCase(), a);
}
}
public ProtocolAdapter getAdapter(String protocol) {
ProtocolAdapter adapter = adapters.get(protocol.toLowerCase());
if (adapter == null) {
throw new IllegalArgumentException("Unsupported protocol: " + protocol);
}
return adapter;
}
}
@@ -0,0 +1,48 @@
package com.water.iot.adapter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.*;
@Slf4j
@Component
public class ModbusAdapter implements ProtocolAdapter {
@Override
public String protocol() { return "Modbus"; }
@Override
public Map<String, Object> parseTelemetry(String deviceSn, byte[] raw) {
// Modbus RTU/TCP 数据帧解析
Map<String, Object> telemetry = new HashMap<>();
telemetry.put("deviceSn", deviceSn);
telemetry.put("timestamp", System.currentTimeMillis());
telemetry.put("raw_hex", bytesToHex(raw));
// 简化: 按寄存器地址映射指标
List<Map<String, Object>> metrics = new ArrayList<>();
Map<String, Object> m = new HashMap<>();
m.put("key", "register_0");
m.put("value", raw.length > 0 ? raw[0] & 0xFF : 0);
metrics.add(m);
telemetry.put("metrics", metrics);
return telemetry;
}
@Override
public byte[] encodeCommand(Map<String, Object> command) {
// Modbus 写寄存器指令
return new byte[]{0x01, 0x06, 0x00, 0x00, 0x00, 0x01, 0x48, 0x0A};
}
@Override
public boolean authenticate(String deviceSn, String credential) {
return true;
}
private String bytesToHex(byte[] bytes) {
StringBuilder sb = new StringBuilder();
for (byte b : bytes) sb.append(String.format("%02X", b));
return sb.toString();
}
}
@@ -0,0 +1,46 @@
package com.water.iot.adapter;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.*;
@Slf4j
@Component
public class MqttAdapter implements ProtocolAdapter {
private final ObjectMapper mapper = new ObjectMapper();
@Override
public String protocol() { return "MQTT"; }
@Override
public Map<String, Object> parseTelemetry(String deviceSn, byte[] raw) {
try {
@SuppressWarnings("unchecked")
Map<String, Object> data = mapper.readValue(raw, Map.class);
Map<String, Object> telemetry = new HashMap<>();
telemetry.put("deviceSn", deviceSn);
telemetry.put("timestamp", System.currentTimeMillis());
// 标准格式: {deviceSn, ts, metrics: [{key, value, unit}]}
telemetry.put("metrics", data.getOrDefault("metrics", data));
return telemetry;
} catch (Exception e) {
log.error("MQTT parse error: {}", e.getMessage());
return null;
}
}
@Override
public byte[] encodeCommand(Map<String, Object> command) {
try { return mapper.writeValueAsBytes(command); }
catch (Exception e) { return null; }
}
@Override
public boolean authenticate(String deviceSn, String credential) {
// TODO: 从设备表查询校验
return true;
}
}
@@ -0,0 +1,22 @@
package com.water.iot.adapter;
import java.util.Map;
/**
* 设备协议适配器接口 — 策略模式
* 每种协议实现此接口,统一处理设备数据
*/
public interface ProtocolAdapter {
/** 支持的协议名 */
String protocol();
/** 将原始数据转为标准遥测格式 */
Map<String, Object> parseTelemetry(String deviceSn, byte[] raw);
/** 将指令转为协议特定的下发格式 */
byte[] encodeCommand(Map<String, Object> command);
/** 设备鉴权 */
boolean authenticate(String deviceSn, String credential);
}
@@ -0,0 +1,49 @@
package com.water.iot.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import java.util.Map;
import java.util.concurrent.TimeUnit;
@Slf4j
@Service
@RequiredArgsConstructor
public class DeviceShadowService {
private final StringRedisTemplate redisTemplate;
private final JdbcTemplate jdbcTemplate;
private final ObjectMapper mapper = new ObjectMapper();
private static final String SHADOW_PREFIX = "iot:shadow:";
private static final long SHADOW_TTL_HOURS = 24;
/** 更新设备上报状态 */
public void updateReported(String deviceSn, Map<String, Object> state) {
try {
String key = SHADOW_PREFIX + deviceSn;
String json = mapper.writeValueAsString(state);
redisTemplate.opsForHash().put(key, "reported", json);
redisTemplate.expire(key, SHADOW_TTL_HOURS, TimeUnit.HOURS);
// 同步更新数据库设备最后上报时间
jdbcTemplate.update("UPDATE iot_device SET last_report_time = NOW() WHERE device_sn = ?", deviceSn);
} catch (JsonProcessingException e) {
log.error("Shadow update error: {}", e.getMessage());
}
}
/** 获取设备影子 */
public Map<Object, Object> getShadow(String deviceSn) {
return redisTemplate.opsForHash().entries(SHADOW_PREFIX + deviceSn);
}
/** 更新期望状态(云端→设备) */
public void updateDesired(String deviceSn, String desiredJson) {
redisTemplate.opsForHash().put(SHADOW_PREFIX + deviceSn, "desired", desiredJson);
}
}
@@ -0,0 +1,34 @@
package com.water.iot.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import java.util.Map;
@Slf4j
@Service
@RequiredArgsConstructor
public class OtaService {
private final JdbcTemplate jdbcTemplate;
private final DeviceShadowService shadowService;
/** 创建 OTA 升级任务 */
public void createUpgrade(Long modelId, String firmwareVersion, String firmwareUrl, String checkMd5) {
jdbcTemplate.update(
"INSERT INTO iot_device_event (device_id, device_sn, event_type, event_data) " +
"SELECT id, device_sn, 'ota', json_build_object('version',?, 'url',?, 'md5',?) " +
"FROM iot_device WHERE model_id = ? AND status = 'online'",
firmwareVersion, firmwareUrl, checkMd5, modelId);
log.info("OTA task created for model {}: version={}", modelId, firmwareVersion);
}
/** 设备查询是否有待升级固件 */
public Map<String, Object> checkUpgrade(String deviceSn, String currentVersion) {
return jdbcTemplate.queryForMap(
"SELECT * FROM iot_device_event WHERE device_sn = ? AND event_type = 'ota' ORDER BY created_at DESC LIMIT 1",
deviceSn);
}
}
@@ -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<Map<String, Object>> createRoute(@RequestBody Map<String, Object> req) {
@SuppressWarnings("unchecked")
List<Map<String, Object>> points = (List<Map<String, Object>>) 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<List<Map<String, Object>>> routes(@RequestParam String area) {
return R.ok(patrolService.getRoutes(area));
}
// ---- 任务 ----
@PostMapping("/task")
public R<Map<String, Object>> createTask(@RequestBody Map<String, Object> 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<List<Map<String, Object>>> todayTasks(@RequestParam Long userId) {
return R.ok(patrolService.getTodayTasks(userId));
}
@PutMapping("/task/{id}/start")
public R<Map<String, Object>> startTask(@PathVariable Long id) {
return R.ok(patrolService.startTask(id));
}
@PutMapping("/task/{id}/complete")
public R<Map<String, Object>> completeTask(@PathVariable Long id, @RequestParam double distance) {
return R.ok(patrolService.completeTask(id, distance));
}
// ---- 巡检记录 ----
@PostMapping("/record")
public R<Map<String, Object>> record(@RequestBody Map<String, Object> req) {
@SuppressWarnings("unchecked")
List<Map<String, Object>> items = (List<Map<String, Object>>) 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<List<Map<String, Object>>> records(@PathVariable Long taskId) {
return R.ok(patrolService.getTaskRecords(taskId));
}
// ---- 问题上报 ----
@PostMapping("/issue/report")
public R<Map<String, Object>> reportIssue(@RequestBody Map<String, Object> req) {
@SuppressWarnings("unchecked")
List<String> photos = (List<String>) 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<Map<String, Object>> stats(@RequestParam String area,
@RequestParam String start,
@RequestParam String end) {
return R.ok(patrolService.getStats(area, LocalDate.parse(start), LocalDate.parse(end)));
}
}
@@ -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<String, Object> createRoute(String routeName, String area, List<Map<String, Object>> 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<Map<String, Object>> getRoutes(String area) {
return jdbc.queryForList("SELECT * FROM patrol_route WHERE area = ? AND status = 1", area);
}
// ========== 任务管理 ==========
public Map<String, Object> 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<Map<String, Object>> 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<String, Object> 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<String, Object> 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<String, Object> recordCheck(Long taskId, int pointSeq, Long deviceId,
List<Map<String, Object>> 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<Map<String, Object>> getTaskRecords(Long taskId) {
return jdbc.queryForList(
"SELECT * FROM patrol_record WHERE task_id = ? ORDER BY point_seq", taskId);
}
// ========== 问题上报(巡检APP) ==========
public Map<String, Object> reportIssue(Long taskId, Long deviceId, String issueType,
String description, List<String> 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<String, Object> getStats(String area, LocalDate start, LocalDate end) {
Map<String, Object> 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;
}
}
@@ -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<Map<String, Object>> overview(@RequestParam(defaultValue = "一体化水厂") String area,
@RequestParam(defaultValue = "admin") String roleType) {
return R.ok(dashboardService.getOverview(area, roleType));
}
// ---- 实时监测 ----
@GetMapping("/monitor/realtime")
public R<List<Map<String, Object>>> 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<List<Map<String, Object>>> cameras(@RequestParam String area) {
return R.ok(videoService.getCameras(area));
}
// ---- 水质 ----
@GetMapping("/quality/chemical/{station}")
public R<Map<String, Object>> chemical(@PathVariable String station) {
return R.ok(wqService.getChemicalMonitoring(station));
}
@PostMapping("/quality/record")
public R<String> addRecord(@RequestBody Map<String, Object> record) {
wqService.addRecord(record);
return R.ok("记录已保存");
}
@GetMapping("/quality/ledger")
public R<List<Map<String, Object>>> ledger(@RequestParam String area, @RequestParam String start, @RequestParam String end) {
return R.ok(wqService.getQualityLedger(area, start, end));
}
// ---- 报警 ----
@GetMapping("/alert/list")
public R<List<Map<String, Object>>> 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<String> confirm(@PathVariable Long id, @RequestParam Long userId) {
alertEngine.confirm(id, userId); return R.ok("已确认");
}
@PostMapping("/alert/{id}/dispatch")
public R<String> dispatch(@PathVariable Long id, @RequestParam Long assigneeId) {
alertEngine.dispatch(id, assigneeId); return R.ok("已派单");
}
// ---- 调度 ----
@GetMapping("/dispatch/duty/today")
public R<List<Map<String, Object>>> todayDuty(@RequestParam String area) {
return R.ok(dispatchService.getTodayDuty(area));
}
@PostMapping("/dispatch/command")
public R<Map<String, Object>> createCommand(@RequestBody Map<String, Object> req) {
@SuppressWarnings("unchecked")
List<Long> targetIds = (List<Long>) 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<Map<String, Object>> issueCommand(@PathVariable String cmdNo) {
return R.ok(dispatchService.issueCommand(cmdNo));
}
@PostMapping("/dispatch/emergency/pipe-burst")
public R<Map<String, Object>> pipeBurst(@RequestBody Map<String, Object> 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<List<Map<String, Object>>> 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<Map<String, Object>> report(@RequestParam String type, @RequestParam String period) {
return R.ok(dataCenterService.generateReport(type, period));
}
@PutMapping("/data/threshold/{ruleId}")
public R<String> updateThreshold(@PathVariable Long ruleId, @RequestBody Map<String, Object> req) {
dataCenterService.updateThreshold(ruleId,
((Number) req.get("threshold")).doubleValue(),
(String) req.get("condition"));
return R.ok("阈值已更新");
}
}
@@ -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<String, Long> lastAlertTime = new ConcurrentHashMap<>();
/** 检查指标是否触发报警 */
public void checkMetric(String deviceSn, String metricKey, double value, String area) {
List<Map<String, Object>> rules = jdbc.queryForList(
"SELECT * FROM alert_rule WHERE metric_key = ? AND enabled = 1", metricKey);
for (Map<String, Object> 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<Map<String, Object>> 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());
}
}
@@ -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<String, Object> getOverview(String area, String roleType) {
Map<String, Object> overview = new LinkedHashMap<>();
// 今日进出水量(从时序库聚合)
try {
Map<String, Object> 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<Map<String, Object>> 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());
}
}
@@ -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<Map<String, Object>> 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<String, Object> generateReport(String reportType, String period) {
Map<String, Object> 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<Map<String, Object>> 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);
}
}
@@ -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<Map<String, Object>> 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<String, Object> 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<String, Object> 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<String, Object> createCommand(String title, String content, String type, String source, String targetType, List<Long> 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<String, Object> 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<String, Object> trackCommand(String cmdNo) {
return jdbc.queryForMap("SELECT * FROM dispatch_command WHERE command_no = ?", cmdNo);
}
public List<Map<String, Object>> 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<String, Object> pipeBurstSimulation(double lng, double lat, String pipeDiameter) {
// 爆管模拟:影响区域分析
Map<String, Object> 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<String, Object> waterQualityIncident(String area, String pollutant) {
Map<String, Object> 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;
}
}
@@ -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<Map<String, Object>> getCameras(String area) {
// Mock: 返回预设视频点位
List<Map<String, Object>> 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<String, Object> 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);
}
}
@@ -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<String, Object> getChemicalMonitoring(String stationName) {
Map<String, Object> 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<Map<String, Object>> 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<Map<String, Object>> 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<String, Object> 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"));
}
}
@@ -0,0 +1,59 @@
package com.water.revenue.controller;
import com.water.common.core.result.R;
import com.water.revenue.service.RemoteReadingService;
import com.water.revenue.service.WorkOrderService;
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("/revenue/ops")
@RequiredArgsConstructor
public class MeterWorkController {
private final RemoteReadingService rrService;
private final WorkOrderService woService;
// ---- 远传集抄 ----
@PostMapping("/reading/batch/{area}")
public R<Map<String, Object>> batchRead(@PathVariable String area) {
return R.ok(rrService.batchRead(area));
}
@GetMapping("/dma/analysis")
public R<Map<String, Object>> dmaAnalysis(@RequestParam String area, @RequestParam String period) {
return R.ok(rrService.dmaAnalysis(area, period));
}
@GetMapping("/meter/large")
public R<List<Map<String, Object>>> largeMeters() {
return R.ok(rrService.largeMeterMonitor());
}
// ---- 工单 ----
@PostMapping("/work-order")
public R<Map<String, Object>> createWO(@RequestBody Map<String, Object> req) {
return R.ok(woService.create(
(String) req.get("title"), (String) req.get("type"), (String) req.get("priority"),
1L, "当前用户", (String) req.get("description"), (String) req.get("area")));
}
@PutMapping("/work-order/assign")
public R<Map<String, Object>> assign(@RequestBody Map<String, Object> req) {
return R.ok(woService.assign(
(String) req.get("woNo"), Long.parseLong(String.valueOf(req.get("assigneeId"))),
(String) req.get("assigneeName")));
}
@PutMapping("/work-order/complete")
public R<Map<String, Object>> complete(@RequestBody Map<String, Object> req) {
@SuppressWarnings("unchecked")
List<String> photos = (List<String>) req.getOrDefault("photos", List.of());
return R.ok(woService.complete((String) req.get("woNo"), (String) req.get("result"), photos));
}
}
@@ -0,0 +1,79 @@
package com.water.revenue.controller;
import com.water.common.core.result.R;
import com.water.revenue.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.math.BigDecimal;
import java.util.*;
@Tag(name = "营业收费")
@RestController
@RequestMapping("/revenue")
@RequiredArgsConstructor
public class RevenueController {
private final RevenueBaseService baseService;
private final InstallService installService;
private final BillingService billingService;
private final MeterService meterService;
// ---- 营收平台 ----
@PostMapping("/auth/sso")
public R<Map<String, Object>> ssoLogin(@RequestBody Map<String, String> req) {
return R.ok(baseService.ssoLogin(req.get("username"), req.get("password"), req.get("appType")));
}
// ---- 报装管理 ----
@PostMapping("/install/pre-apply")
public R<Map<String, Object>> preApply(@RequestBody Map<String, String> req) {
return R.ok(installService.preApply(req.get("name"), req.get("phone"),
req.get("area"), req.get("address"), req.get("customerType"), req.get("caliber")));
}
@GetMapping("/install/progress/{appNo}")
public R<Map<String, Object>> progress(@PathVariable String appNo) {
return R.ok(installService.getProgress(appNo));
}
// ---- 营业收费 ----
@PostMapping("/billing/generate")
public R<Map<String, Object>> generateBill(@RequestParam Long readingId) {
return R.ok(billingService.generateBill(readingId));
}
@PostMapping("/billing/pay")
public R<Map<String, Object>> pay(@RequestBody Map<String, Object> req) {
return R.ok(billingService.pay(
Long.parseLong(String.valueOf(req.get("billId"))),
(String) req.get("payMethod"),
(String) req.get("payChannel"),
new BigDecimal(String.valueOf(req.get("amount")))));
}
// ---- 表务管理 ----
@PostMapping("/meter/stock-in")
public R<String> stockIn(@RequestBody Map<String, String> req) {
meterService.stockIn(req.get("meterNo"), req.get("caliber"),
req.get("meterType"), req.get("manufacturer"), Integer.parseInt(req.get("quantity")));
return R.ok("入库成功");
}
@PostMapping("/meter/replace")
public R<String> replace(@RequestBody Map<String, Object> req) {
meterService.replace(
Long.parseLong(String.valueOf(req.get("oldMeterId"))),
Long.parseLong(String.valueOf(req.get("newMeterId"))),
new BigDecimal(String.valueOf(req.get("oldReading"))),
(String) req.get("remark"));
return R.ok("换表成功");
}
@GetMapping("/meter/lifecycle/{meterId}")
public R<List<Map<String, Object>>> lifecycle(@PathVariable Long meterId) {
return R.ok(meterService.getLifecycle(meterId));
}
}
@@ -0,0 +1,65 @@
package com.water.revenue.controller;
import com.water.common.core.result.R;
import com.water.revenue.service.CustomerServiceCenter;
import com.water.revenue.service.WechatService;
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.math.BigDecimal;
import java.util.*;
@Tag(name = "客服热线 & 微信网厅")
@RestController
@RequestMapping("/revenue/customer")
@RequiredArgsConstructor
public class WechatController {
private final CustomerServiceCenter csc;
private final WechatService wechatService;
// ---- 客服 ----
@GetMapping("/bills/query")
public R<List<Map<String, Object>>> queryBills(@RequestParam String phoneOrNo) {
return R.ok(csc.queryBills(phoneOrNo));
}
@GetMapping("/knowledge/search")
public R<List<Map<String, Object>>> searchKnowledge(@RequestParam String keyword) {
return R.ok(csc.searchKnowledge(keyword));
}
@GetMapping("/notices/{type}")
public R<List<Map<String, Object>>> notices(@PathVariable String type) {
return R.ok(csc.getNotices(type));
}
@GetMapping("/kpi")
public R<Map<String, Object>> kpi() { return R.ok(csc.getKpi()); }
// ---- 微信网厅 ----
@PostMapping("/wechat/bind")
public R<Map<String, Object>> bind(@RequestBody Map<String, String> req) {
return R.ok(wechatService.bindUser(req.get("openId"), req.get("customerNo"), req.get("phone")));
}
@PostMapping("/wechat/pay/prepay")
public R<Map<String, Object>> prepay(@RequestBody Map<String, String> req) {
return R.ok(wechatService.wechatPayPrepay(
req.get("customerNo"), req.get("billPeriod"),
new BigDecimal(req.get("amount"))));
}
@GetMapping("/wechat/ai-answer")
public R<String> aiAnswer(@RequestParam String question) {
return R.ok(wechatService.aiAnswer(question));
}
@PostMapping("/wechat/notice/publish")
public R<String> publishNotice(@RequestBody Map<String, String> req) {
wechatService.publishNotice(req.get("type"), req.get("title"), req.get("content"));
return R.ok("公告已发布");
}
}
@@ -0,0 +1,96 @@
package com.water.revenue.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.LocalDate;
import java.time.YearMonth;
import java.time.format.DateTimeFormatter;
import java.util.*;
@Slf4j
@Service
@RequiredArgsConstructor
public class BillingService {
private final JdbcTemplate jdbcTemplate;
/** 根据抄表记录生成账单 */
@Transactional
public Map<String, Object> generateBill(Long readingId) {
Map<String, Object> reading = jdbcTemplate.queryForMap(
"SELECT r.*, m.customer_id, m.caliber, c.customer_type FROM rev_reading r " +
"JOIN rev_meter m ON r.meter_id = m.id " +
"JOIN rev_customer c ON m.customer_id = c.id " +
"WHERE r.id = ?", readingId);
BigDecimal consumption = (BigDecimal) reading.get("consumption");
String customerType = (String) reading.get("customer_type");
String period = (String) reading.get("reading_period");
// 查询阶梯水价
List<Map<String, Object>> prices = jdbcTemplate.queryForList(
"SELECT * FROM rev_water_price WHERE customer_type = ? AND effective_date <= CURRENT_DATE ORDER BY tier_no",
customerType);
BigDecimal waterFee = BigDecimal.ZERO;
BigDecimal remaining = consumption;
for (Map<String, Object> p : prices) {
BigDecimal rangeEnd = p.get("range_end") != null ? (BigDecimal) p.get("range_end") : BigDecimal.valueOf(99999);
BigDecimal price = (BigDecimal) p.get("water_price");
BigDecimal tierUsage = remaining.min(rangeEnd);
waterFee = waterFee.add(tierUsage.multiply(price));
remaining = remaining.subtract(tierUsage);
if (remaining.compareTo(BigDecimal.ZERO) <= 0) break;
}
// 污水处理费(用水量的80%)
BigDecimal sewageFee = consumption.multiply(((BigDecimal) prices.get(0).getOrDefault("sewage_price", BigDecimal.ZERO)))
.multiply(BigDecimal.valueOf(0.8));
BigDecimal totalFee = waterFee.add(sewageFee);
String billNo = "BILL-" + System.currentTimeMillis();
Date dueDate = java.sql.Date.valueOf(LocalDate.now().plusDays(30));
jdbcTemplate.update(
"INSERT INTO rev_bill (bill_no, customer_id, meter_id, reading_id, bill_period, prev_reading, curr_reading, consumption, water_fee, sewage_fee, total_fee, status, due_date) " +
"VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)",
billNo, reading.get("customer_id"), reading.get("meter_id"), readingId,
period, reading.get("prev_reading"), reading.get("curr_reading"),
consumption, waterFee, sewageFee, totalFee, "pending", dueDate);
log.info("Bill generated: {} total={}", billNo, totalFee);
return Map.of("billNo", billNo, "totalFee", totalFee, "waterFee", waterFee, "sewageFee", sewageFee);
}
/** 缴费 */
@Transactional
public Map<String, Object> pay(long billId, String payMethod, String payChannel, BigDecimal amount) {
String paymentNo = "PAY-" + System.currentTimeMillis();
jdbcTemplate.update(
"INSERT INTO rev_payment (bill_id, customer_id, payment_no, amount, pay_method, pay_channel) " +
"SELECT ?, customer_id, ?, ?, ?, ? FROM rev_bill WHERE id = ?",
billId, paymentNo, amount, payMethod, payChannel, billId);
jdbcTemplate.update(
"UPDATE rev_bill SET paid_fee = paid_fee + ?, status = CASE WHEN paid_fee >= total_fee THEN 'paid' ELSE 'partial' END, paid_at = NOW() WHERE id = ?",
amount, billId);
return Map.of("paymentNo", paymentNo, "amount", amount, "status", "success");
}
/** 欠费统计 */
public List<Map<String, Object>> getOverdueBills(String area) {
return jdbcTemplate.queryForList(
"SELECT b.*, c.customer_name, c.phone FROM rev_bill b " +
"JOIN rev_customer c ON b.customer_id = c.id " +
"WHERE b.status IN ('pending','partial') AND c.area = ? AND b.due_date < CURRENT_DATE",
area);
}
}
@@ -0,0 +1,52 @@
package com.water.revenue.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 CustomerServiceCenter {
private final JdbcTemplate jdbcTemplate;
/** 水费查询(语音/在线) */
public List<Map<String, Object>> queryBills(String phoneOrCustomerNo) {
return jdbcTemplate.queryForList(
"SELECT b.*, c.customer_name, c.phone FROM rev_bill b " +
"JOIN rev_customer c ON b.customer_id = c.id " +
"WHERE c.phone = ? OR c.customer_no = ? ORDER BY b.bill_period DESC LIMIT 12",
phoneOrCustomerNo, phoneOrCustomerNo);
}
/** 知识库管理 */
public List<Map<String, Object>> searchKnowledge(String keyword) {
return jdbcTemplate.queryForList(
"SELECT dict_label, dict_value FROM sys_dict_data WHERE dict_type_id = " +
"(SELECT id FROM sys_dict_type WHERE dict_key = 'knowledge_base') " +
"AND dict_label LIKE ?", "%" + keyword + "%");
}
/** 公告板 */
public List<Map<String, Object>> getNotices(String noticeType) {
return jdbcTemplate.queryForList(
"SELECT dict_label, dict_value, created_at FROM sys_dict_data WHERE dict_type_id = " +
"(SELECT id FROM sys_dict_type WHERE dict_key = ?) ORDER BY created_at DESC LIMIT 10",
"notice_" + noticeType); // notice_water_stop, notice_water_quality, etc.
}
/** KPI 指标 */
public Map<String, Object> getKpi() {
return jdbcTemplate.queryForMap("""
SELECT
(SELECT COUNT(*) FROM rev_bill WHERE status = 'pending' AND created_at > CURRENT_DATE - 30) AS pending_bills,
(SELECT COUNT(*) FROM rev_install WHERE status IN ('pre_apply','engineering')) AS pending_installs,
(SELECT ROUND(AVG(EXTRACT(EPOCH FROM (updated_at - created_at))/3600)::numeric, 2)
FROM rev_install WHERE status = 'completed' AND created_at > CURRENT_DATE - 30) AS avg_install_hours
""");
}
}
@@ -0,0 +1,56 @@
package com.water.revenue.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 InstallService {
private final JdbcTemplate jdbcTemplate;
/** 预受理申请 */
public Map<String, Object> preApply(String name, String phone, String area, String address, String customerType, String caliber) {
String appNo = "INS-" + System.currentTimeMillis();
jdbcTemplate.update(
"INSERT INTO rev_install (application_no, applicant_name, applicant_phone, area, address, customer_type, caliber, status) VALUES (?,?,?,?,?,?,?,?)",
appNo, name, phone, area, address, customerType, caliber, "pre_apply");
return Map.of("applicationNo", appNo, "status", "pre_apply");
}
/** 工程申请 */
public Map<String, Object> engineeringApply(String appNo, Map<String, Object> engData) {
jdbcTemplate.update(
"UPDATE rev_install SET status = 'engineering', updated_at = NOW() WHERE application_no = ?",
appNo);
return Map.of("applicationNo", appNo, "status", "engineering");
}
/** 派单到施工 */
public Map<String, Object> assignTask(String appNo, Long assigneeId) {
jdbcTemplate.update(
"UPDATE rev_install SET status = 'pending_review', updated_at = NOW() WHERE application_no = ?",
appNo);
return Map.of("applicationNo", appNo, "status", "pending_review", "assigneeId", assigneeId);
}
/** 查询报装进度 */
public Map<String, Object> getProgress(String appNo) {
return jdbcTemplate.queryForMap(
"SELECT application_no, applicant_name, applicant_phone, area, address, customer_type, caliber, status, created_at, updated_at FROM rev_install WHERE application_no = ?",
appNo);
}
/** 报装统计报表 */
public List<Map<String, Object>> getStatsReport(String area, LocalDate start, LocalDate end) {
return jdbcTemplate.queryForList(
"SELECT area, customer_type, status, COUNT(*) as count FROM rev_install WHERE created_at BETWEEN ? AND ? GROUP BY area, customer_type, status",
start, end);
}
}
@@ -0,0 +1,74 @@
package com.water.revenue.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.math.BigDecimal;
import java.util.*;
@Slf4j
@Service
@RequiredArgsConstructor
public class MeterService {
private final JdbcTemplate jdbcTemplate;
/** 水表入库 */
public void stockIn(String meterNo, String caliber, String meterType, String manufacturer, int quantity) {
for (int i = 0; i < quantity; i++) {
String no = meterNo + "-" + (i + 1);
jdbcTemplate.update(
"INSERT INTO rev_meter (meter_no, caliber, meter_type, manufacturer, status) VALUES (?,?,?,?,?)",
no, caliber, meterType, manufacturer, "warehouse");
}
log.info("Meter stock in: {} x{}", meterNo, quantity);
}
/** 水表出库安装 */
@Transactional
public void install(Long meterId, Long customerId, String address, BigDecimal initialReading) {
jdbcTemplate.update(
"UPDATE rev_meter SET customer_id = ?, install_address = ?, initial_reading = ?, current_reading = ?, status = 'active', install_date = CURRENT_DATE WHERE id = ?",
customerId, address, initialReading, initialReading, meterId);
jdbcTemplate.update(
"INSERT INTO rev_meter_log (meter_id, operation_type, old_reading, new_reading, remark) VALUES (?,?,?,?,?)",
meterId, "install", BigDecimal.ZERO, initialReading, "新表安装");
}
/** 故障换表 */
@Transactional
public void replace(Long oldMeterId, Long newMeterId, BigDecimal oldReading, String remark) {
// 旧表拆除
jdbcTemplate.update("UPDATE rev_meter SET status = 'dismantled', current_reading = ? WHERE id = ?", oldReading, oldMeterId);
jdbcTemplate.update(
"INSERT INTO rev_meter_log (meter_id, operation_type, old_reading, remark) VALUES (?,?,?,?)",
oldMeterId, "dismantle", oldReading, remark);
// 新表安装(继承客户信息)
Map<String, Object> oldMeter = jdbcTemplate.queryForMap("SELECT customer_id, install_address FROM rev_meter WHERE id = ?", oldMeterId);
jdbcTemplate.update(
"UPDATE rev_meter SET customer_id = ?, install_address = ?, status = 'active', install_date = CURRENT_DATE WHERE id = ?",
oldMeter.get("customer_id"), oldMeter.get("install_address"), newMeterId);
jdbcTemplate.update(
"INSERT INTO rev_meter_log (meter_id, operation_type, new_meter_no, remark) VALUES (?,?,?,?)",
newMeterId, "change", String.valueOf(newMeterId), "替换旧表 #" + oldMeterId);
}
/** 水表报废 */
public void scrap(Long meterId, String reason) {
jdbcTemplate.update("UPDATE rev_meter SET status = 'scrapped' WHERE id = ?", meterId);
jdbcTemplate.update(
"INSERT INTO rev_meter_log (meter_id, operation_type, remark) VALUES (?,?,?)",
meterId, "scrap", reason);
}
/** 查询水表生命周期记录 */
public List<Map<String, Object>> getLifecycle(Long meterId) {
return jdbcTemplate.queryForList(
"SELECT * FROM rev_meter_log WHERE meter_id = ? ORDER BY created_at",
meterId);
}
}
@@ -0,0 +1,78 @@
package com.water.revenue.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.*;
@Slf4j
@Service
@RequiredArgsConstructor
public class RemoteReadingService {
private final JdbcTemplate jdbcTemplate;
/** 远传集抄:批量采集 */
public Map<String, Object> batchRead(String area) {
List<Map<String, Object>> meters = jdbcTemplate.queryForList(
"SELECT rm.id, rm.meter_no, rm.current_reading, i.device_sn " +
"FROM rev_meter rm LEFT JOIN iot_device i ON rm.device_id = i.id " +
"JOIN rev_customer c ON rm.customer_id = c.id " +
"WHERE rm.status = 'active' AND c.area = ?", area);
int success = 0, failed = 0;
String period = java.time.YearMonth.now().format(java.time.format.DateTimeFormatter.ofPattern("yyyy-MM"));
for (Map<String, Object> m : meters) {
try {
String deviceSn = (String) m.get("device_sn");
// 从 IoT 平台获取实时读数(mock: 随机增量)
BigDecimal prev = m.get("current_reading") != null ? (BigDecimal) m.get("current_reading") : BigDecimal.ZERO;
BigDecimal curr = prev.add(BigDecimal.valueOf(new Random().nextDouble() * 50));
BigDecimal consumption = curr.subtract(prev);
if (consumption.compareTo(BigDecimal.ZERO) < 0) consumption = BigDecimal.ZERO;
jdbcTemplate.update(
"INSERT INTO rev_reading (meter_id, reading_date, reading_period, prev_reading, curr_reading, consumption, read_type) " +
"VALUES (?, CURRENT_DATE, ?, ?, ?, ?, 'remote')",
m.get("id"), period, prev, curr, consumption);
jdbcTemplate.update("UPDATE rev_meter SET current_reading = ? WHERE id = ?", curr, m.get("id"));
success++;
} catch (Exception e) {
failed++;
log.warn("Read failed for meter {}: {}", m.get("meter_no"), e.getMessage());
}
}
log.info("Batch read: area={} success={} failed={}", area, success, failed);
return Map.of("area", area, "success", success, "failed", failed, "period", period);
}
/** DMA 分区漏损分析 */
public Map<String, Object> dmaAnalysis(String area, String dateStr) {
List<Map<String, Object>> result = jdbcTemplate.queryForList(
"SELECT area, SUM(consumption) as total_consumption, COUNT(DISTINCT rm.id) as meter_count " +
"FROM rev_reading rr JOIN rev_meter rm ON rr.meter_id = rm.id " +
"JOIN rev_customer c ON rm.customer_id = c.id " +
"WHERE c.area = ? AND rr.reading_period = ? GROUP BY area",
area, dateStr);
// 漏损率 = 1 - (售水量/供水量)
Map<String, Object> dma = new HashMap<>(result.isEmpty() ? Map.of() : result.get(0));
dma.put("supplyEstimate", 1000); // TODO: 从水厂出水量获取
dma.put("leakRate", "分析中");
return dma;
}
/** 大表监控 (DN80+) */
public List<Map<String, Object>> largeMeterMonitor() {
return jdbcTemplate.queryForList(
"SELECT rm.*, c.customer_name, c.area, i.device_sn, i.status as device_status " +
"FROM rev_meter rm JOIN rev_customer c ON rm.customer_id = c.id " +
"LEFT JOIN iot_device i ON rm.device_id = i.id " +
"WHERE rm.caliber IN ('DN80','DN100','DN150','DN200','DN300','DN400') AND rm.status = 'active'");
}
}
@@ -0,0 +1,49 @@
package com.water.revenue.service;
import cn.dev33.satoken.secure.BCrypt;
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 RevenueBaseService {
private final JdbcTemplate jdbcTemplate;
// ========== SSO 单点登录 ==========
public Map<String, Object> ssoLogin(String username, String password, String appType) {
// 从统一用户表验证
Map<String, Object> user = jdbcTemplate.queryForMap(
"SELECT id, username, real_name, phone, status FROM sys_user WHERE username = ? AND status = 1",
username);
if (user == null) throw new RuntimeException("用户不存在");
// 生成 SSO Token
String ssoToken = UUID.randomUUID().toString();
jdbcTemplate.update(
"INSERT INTO sys_oper_log (user_id, username, module, operation, request_url) VALUES (?,?,?,?,?)",
user.get("id"), username, "revenue", "sso_login", "/revenue/auth/sso");
user.put("ssoToken", ssoToken);
user.put("appType", appType);
return user;
}
// ========== 应用接入管理 ==========
public void registerApp(String appName, String appKey, String appSecret, String redirectUri) {
jdbcTemplate.update(
"INSERT INTO sys_dict_data (dict_type_id, dict_label, dict_value) " +
"SELECT id, ?, ? FROM sys_dict_type WHERE dict_key = 'app_config'",
appName, appKey + ":" + appSecret + ":" + redirectUri);
}
// ========== 运维审计 ==========
public void auditLog(Long userId, String action, String target, String detail) {
jdbcTemplate.update(
"INSERT INTO sys_oper_log (user_id, module, operation, request_url, request_params) VALUES (?,?,?,?,?)",
userId, "revenue_audit", action, target, detail);
}
}
@@ -0,0 +1,68 @@
package com.water.revenue.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import java.math.BigDecimal;
import java.util.*;
@Slf4j
@Service
@RequiredArgsConstructor
public class WechatService {
private final JdbcTemplate jdbcTemplate;
/** 微信用户绑定 */
public Map<String, Object> bindUser(String wechatOpenId, String customerNo, String phone) {
jdbcTemplate.update(
"UPDATE rev_customer SET phone = COALESCE(phone, ?) WHERE customer_no = ?",
phone, customerNo);
// 存储微信绑定关系
jdbcTemplate.update(
"INSERT INTO sys_dict_data (dict_type_id, dict_label, dict_value) " +
"SELECT id, ?, ? FROM sys_dict_type WHERE dict_key = 'wechat_bind' " +
"ON CONFLICT DO NOTHING",
wechatOpenId, customerNo);
return Map.of("openId", wechatOpenId, "customerNo", customerNo, "status", "bound");
}
/** 在线缴费(微信支付预下单) */
public Map<String, Object> wechatPayPrepay(String customerNo, String billPeriod, BigDecimal amount) {
String orderNo = "WXPAY-" + System.currentTimeMillis();
Map<String, Object> bill = jdbcTemplate.queryForMap(
"SELECT id, total_fee FROM rev_bill WHERE customer_id = " +
"(SELECT id FROM rev_customer WHERE customer_no = ?) AND bill_period = ? AND status IN ('pending','partial')",
customerNo, billPeriod);
// 微信支付统一下单(mock)
log.info("WeChat Pay prepay: order={} bill={} amount={}", orderNo, bill.get("id"), amount);
return Map.of("orderNo", orderNo, "prepayId", "wx" + orderNo, "amount", amount);
}
/** AI 客服问答 */
public String aiAnswer(String question) {
// 简易关键词匹配
Map<String, String> qa = new LinkedHashMap<>();
qa.put("水费", "您可以发送户号查询水费账单,或通过在线缴费功能直接支付。");
qa.put("停水", "请查看停水公告了解最新停水计划。如有紧急停水,请拨打客服热线。");
qa.put("报装", "新装水表可通过网上营业厅-业务办理-报装申请提交,我们会安排现场踏勘。");
qa.put("水质", "水质报告每月更新,详见水质公告。如发现水质异常请立即联系我们。");
qa.put("发票", "缴费后可在电子发票中申请开具电子发票,发送到您的微信或邮箱。");
qa.put("过户", "房屋买卖后请携带房产证和身份证前往营业厅办理过户手续。");
for (Map.Entry<String, String> e : qa.entrySet()) {
if (question.contains(e.getKey())) return e.getValue();
}
return "您好!我是智慧水务AI客服,您可以问我关于水费、报装、停水、水质、发票等问题。如需人工服务请转接客服热线。";
}
/** 后台公告管理 */
public void publishNotice(String type, String title, String content) {
jdbcTemplate.update(
"INSERT INTO sys_dict_data (dict_type_id, dict_label, dict_value) " +
"SELECT id, ?, ? FROM sys_dict_type WHERE dict_key = ?",
title, content, "notice_" + type);
}
}
@@ -0,0 +1,49 @@
package com.water.revenue.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 WorkOrderService {
private final JdbcTemplate jdbcTemplate;
/** 创建工单 */
public Map<String, Object> create(String title, String type, String priority, Long reporterId,
String reporterName, String description, String area) {
String woNo = "WO-" + System.currentTimeMillis();
jdbcTemplate.update(
"INSERT INTO patrol_task (task_name, task_date, status) VALUES (?, CURRENT_DATE, 'pending')",
title + "[" + woNo + "]");
log.info("WorkOrder created: {} type={} priority={}", woNo, type, priority);
return Map.of("woNo", woNo, "title", title, "status", "pending", "area", area);
}
/** 工单分派 */
public Map<String, Object> assign(String woNo, Long assigneeId, String assigneeName) {
jdbcTemplate.update(
"UPDATE patrol_task SET assignee_id = ?, status = 'in_progress' WHERE task_name LIKE ?",
assigneeId, "%" + woNo + "%");
return Map.of("woNo", woNo, "assigneeId", assigneeId, "status", "assigned");
}
/** 工单处理完成 */
public Map<String, Object> complete(String woNo, String result, List<String> photoUrls) {
jdbcTemplate.update(
"UPDATE patrol_task SET status = 'completed', actual_end = NOW() WHERE task_name LIKE ?",
"%" + woNo + "%");
return Map.of("woNo", woNo, "status", "completed", "photos", photoUrls);
}
/** 工单统计 */
public Map<String, Object> stats(String area) {
return jdbcTemplate.queryForMap(
"SELECT status, COUNT(*) as count FROM patrol_task WHERE task_date >= CURRENT_DATE - 30 GROUP BY status");
}
}