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 标准化管道演示
This commit is contained in:
bot_pm
2026-06-14 13:25:02 +08:00
parent 28dcea5fb6
commit 919f75cf9e
3 changed files with 219 additions and 0 deletions
@@ -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);
}
}