diff --git a/wm-data-engine/src/main/java/com/water/data_engine/controller/DataController.java b/wm-data-engine/src/main/java/com/water/data_engine/controller/DataController.java new file mode 100644 index 00000000..20e4b8bf --- /dev/null +++ b/wm-data-engine/src/main/java/com/water/data_engine/controller/DataController.java @@ -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 ingest(@RequestBody Map req) { + String sourceType = (String) req.get("sourceType"); + String sourceId = (String) req.get("sourceId"); + @SuppressWarnings("unchecked") + Map data = (Map) req.get("data"); + collectService.ingest(sourceType, sourceId, data); + return R.ok("数据已接入"); + } + + @Operation(summary = "批量接入") + @PostMapping("/ingest/batch") + public R batchIngest(@RequestBody List> batch) { + collectService.batchIngest(batch); + return R.ok("批量接入完成"); + } + + @Operation(summary = "数据标准化+清洗+质控(管道演示)") + @PostMapping("/pipeline") + public R> pipeline(@RequestBody Map raw) { + Map std = governanceService.standardize(raw); + Map cleaned = governanceService.clean(std); + Map result = governanceService.qualityCheck(cleaned); + return R.ok(result); + } +} diff --git a/wm-data-engine/src/main/java/com/water/data_engine/service/DataCollectService.java b/wm-data-engine/src/main/java/com/water/data_engine/service/DataCollectService.java new file mode 100644 index 00000000..2b7b640c --- /dev/null +++ b/wm-data-engine/src/main/java/com/water/data_engine/service/DataCollectService.java @@ -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 kafkaTemplate; + private final JdbcTemplate jdbcTemplate; + private final ObjectMapper mapper = new ObjectMapper(); + + /** 数据汇聚入口:接收各来源数据,统一写入 Kafka */ + public void ingest(String sourceType, String sourceId, Map rawData) { + try { + Map 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 envelope = mapper.readValue(message, Map.class); + @SuppressWarnings("unchecked") + Map data = (Map) envelope.get("data"); + + String deviceSn = (String) data.getOrDefault("deviceSn", "unknown"); + @SuppressWarnings("unchecked") + List> metrics = (List>) data.getOrDefault("metrics", List.of()); + + for (Map 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> batchData) { + for (Map data : batchData) { + String sourceType = (String) data.getOrDefault("sourceType", "batch"); + String sourceId = (String) data.getOrDefault("sourceId", UUID.randomUUID().toString()); + @SuppressWarnings("unchecked") + Map rawData = (Map) data.getOrDefault("data", new HashMap<>()); + ingest(sourceType, sourceId, rawData); + } + } +} diff --git a/wm-data-engine/src/main/java/com/water/data_engine/service/DataGovernanceService.java b/wm-data-engine/src/main/java/com/water/data_engine/service/DataGovernanceService.java new file mode 100644 index 00000000..6d716aa4 --- /dev/null +++ b/wm-data-engine/src/main/java/com/water/data_engine/service/DataGovernanceService.java @@ -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 standardize(Map raw) { + Map std = new LinkedHashMap<>(); + // 水利行业标准字段映射 + Map 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 entry : raw.entrySet()) { + String key = standardFields.getOrDefault(entry.getKey(), entry.getKey()); + std.put(key, entry.getValue()); + } + std.put("standardized", true); + return std; + } + + /** 数据清洗:缺失值填充、异常值检测 */ + public Map clean(Map data) { + Map 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 qualityCheck(Map data) { + Map result = new LinkedHashMap<>(data); + int score = 100; + List 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); + } +}