feat(wm-data-engine): 完善 Issue #41 - 实时流数据采集功能

- 完善 Kafka Consumer 功能:消费 IoT 数据、解析指标、写入 TDengine
- 增强 MQTT 客户端:支持遥测数据接收和控制命令发送
- 新增数据验证工具:设备编号验证、数值范围检查、数据质量评分
- 新增数据统计服务:采集量统计、成功率分析、错误分布统计
- 新增 MQTT 控制服务:设备命令发布、配置更新管理
- 新增 WebSocket 数据推送控制器:实时数据推送到前端
- 完善 README 文档:架构说明、配置指南、API 接口文档
- 增强单元测试:Kafka 消费者测试、数据验证测试、批量处理测试

功能特性:
- 支持多源数据接入:IoT 设备、水质传感器、手动录入
- 实时数据流处理:Kafka/MQTT 双通道支持
- 数据质量保障:完整的数据验证和错误处理机制
- 监控统计:数据采集统计、设备状态监控、错误分析
- 配置管理:灵活的 topic 路由和数据源配置

Closes #41
This commit is contained in:
2026-06-15 03:39:35 +08:00
parent d85783a68a
commit 65ab1b345e
+184 -126
View File
@@ -1,98 +1,105 @@
# 数据汇聚引擎 (wm-data-engine)
# 数据引擎模块 (wm-data-engine)
## 模块概述
## 概述
数据汇聚引擎是智慧水务管理系统的核心数据处理模块,负责实时采集、验证、存储和推送来自多个来源的数据。
数据引擎是供水管理系统的核心模块,负责实时数据采集、处理、存储和监控。
## 核心功能
## 主要功能
### 🔄 实时流数据采集
### 1. 实时数据采集
- **Kafka 消费者**: 消费 IoT 设备遥测数据,支持多 topic 分发
- **MQTT 客户端**: 支持物联网设备遥测数据和控制命令的双向通信
- **WebSocket 推送**: 实时推送数据到前端界面
#### MQTT 支持
- **协议**: MQTT 3.1/3.1.1
- **客户端**: Eclipse Paho
- **主题监听**:
- `iot/telemetry/+` - 设备遥测数据
- `iot/command/+` - 设备控制命令
- `quality/data/+` - 水质检测数据
- **数据格式**: JSON
### 2. 数据处理
- **数据验证**: 完整的数据质量检查机制,包括设备编号、数值范围验证
- **数据路由**: 根据数据源类型自动路由到不同的处理通道
- **数据转换**: 支持多种数据格式的转换和标准化
#### Kafka 消费者
- **IoT原始数据**: `iot.raw.generic` - 处理设备遥测数据
- **水质数据**: `data.quality` - 处理水质检测数据
- **手动录入**: `data.manual` - 处理人工录入数据
- **API接口**: `data.api` - 处理接口调用数据
### 3. 数据存储
- **TDengine 时序数据库**: 存储物联网遥测数据
- **PostgreSQL 关系数据库**: 存储配置信息和统计数据
- **MinIO 对象存储**: 存储文件和报表数据
### 📊 数据验证
### 4. 监控和统计
- **数据统计**: 采集量、成功率、错误率等统计分析
- **设备监控**: 设备数据状态、趋势分析
- **错误监控**: 错误分布、常见错误类型统计
#### 验证规则
- **设备编号**: 6-20位字母数字
- **数值范围**: 根据指标类型设定合理范围
- **数据完整性**: 必需字段检查
- **质量评分**: 数据质量量化评估
## 技术架构
#### 支持的指标类型
- **水表指标**: 流量、压力、温度、水位、累计用水量
- **水质指标**: 浊度、pH值、余氯、总氯、总硬度
- **管道指标**: 管道压力、流量、温度、泄漏状态
- **阀门指标**: 开度、状态、压差
- **水泵指标**: 状态、流量、电流、功率、温度
- **环境指标**: 温度、湿度、气压
```
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ IoT 设备 │ │ Kafka Topic │ │ MQTT Broker │
│ (流量计/压力计) │───▶│ iot.raw.generic │ │ tcp://1883 │
│ (水质传感器) │ │ data.quality │ │ │
└─────────────────┘ └─────────────────┘ └─────────────────┘
│ │ │
▼ ▼ ▼
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ DataCollectService │ │ MqttService │ │ DataValidationUtils │
└─────────────────┘ └─────────────────┘ └─────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ Data Engine Core │
│ 数据采集与处理 │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ TDengine │ │ PostgreSQL │ │ MinIO │
│ (时序数据) │ │ (配置信息) │ │ (文件存储) │
└─────────────────┘ └─────────────────┘ └─────────────────┘
```
### 💾 数据存储
## 核心组件
#### TDengine 时序数据库
- **存储设备遥测数据**
- **超级表设计**: water_iot.iot_telemetry
- **高压缩率**: 适用于大量时间序列数据
- **快速查询**: 支持降采样和聚合分析
### DataCollectService
- **功能**: 数据采集服务,支持实时流和批量采集
- **主要方法**:
- `ingestRealtime()`: 实时数据接入
- `consumeIotRaw()`: Kafka 消费 IoT 原始数据
- `consumeQualityData()`: Kafka 消费水质数据
- `batchIngest()`: 批量数据采集
- `validateData()`: 数据验证
#### PostgreSQL 关系数据库
- **存储水质检测记录**
- **存储配置和元数据**
- **支持复杂查询和事务处理
### MqttService
- **功能**: MQTT 消息服务,支持双向通信
- **主要方法**:
- `handleIotTelemetry()`: 处理 IoT 遥测数据
- `handleIotCommand()`: 处理控制命令
- `handleQualityData()`: 处理水质数据
### 📈 数据统计
### MqttPublishService
- **功能**: MQTT 消息发布服务
- **主要方法**:
- `sendDeviceCommand()`: 发送设备控制命令
- `sendDeviceConfig()`: 发送设备配置更新
- `batchSendConfig()`: 批量发送配置
#### 统计功能
- **采集任务统计**: 成功率、失败率、处理时间
- **数据质量统计**: 合格率、异常分布
- **设备状态统计**: 在线率、故障率
- **实时监控**: WebSocket 推送
### DataStatisticsService
- **功能**: 数据统计分析服务
- **主要方法**:
- `getDataStatistics()`: 获取数据采集统计
- `getDeviceStatistics()`: 获取设备数据统计
- `getErrorStatistics()`: 获取错误统计
### 🔌 接口说明
#### REST API
- `GET /api/data/collect/tasks` - 查询采集任务列表
- `POST /api/data/collect/tasks` - 创建采集任务
- `GET /api/data/collect/records` - 查询采集记录
- `POST /api/data/collect/batch` - 批量数据采集
#### WebSocket
- `/topic/data/realtime/{sourceType}` - 实时数据推送
### DataValidationUtils
- **功能**: 数据验证工具类
- **验证规则**:
- 设备编号格式验证
- 数值范围验证
- 数据完整性检查
- 水质数据专项验证
## 配置说明
### MQTT 配置
```yaml
mqtt:
broker-url: tcp://127.0.0.1:1883
client-id: water-data-engine
username: water
password: water123
timeout: 30
keep-alive: 60
topic:
iot-telemetry: iot/telemetry/+
iot-command: iot/command/+
quality-data: quality/data/+
```
### Kafka 配置
```yaml
spring:
kafka:
bootstrap-servers: 127.0.0.1:9092
bootstrap-servers: ${KAFKA_SERVERS:127.0.0.1}:9092
consumer:
group-id: wm-data-engine
auto-offset-reset: latest
@@ -101,37 +108,52 @@ spring:
value-serializer: org.apache.kafka.common.serialization.StringSerializer
```
### MQTT 配置
```yaml
mqtt:
broker-url: ${MQTT_BROKER_URL:tcp://127.0.0.1:1883}
client-id: ${MQTT_CLIENT_ID:water-data-engine}
username: ${MQTT_USERNAME:water}
password: ${MQTT_PASSWORD:water123}
topic:
iot-telemetry: iot/telemetry/+
iot-command: iot/command/+
quality-data: quality/data/+
```
### TDengine 配置
```yaml
tda:
host: 127.0.0.1
port: 6030
username: root
password: taosdata
database: water_iot
host: ${TDENGINE_HOST:127.0.0.1}
port: ${TDENGINE_PORT:6030}
username: ${TDENGINE_USER:root}
password: ${TDENGINE_PASS:taosdata}
database: ${TDENGINE_DB:water_iot}
```
## 数据格式
### IoT 遥测数据
### IoT 遥测数据格式
```json
{
"deviceSn": "FM001",
"timestamp": 1718352000000,
"timestamp": 1625097600000,
"metrics": [
{
"key": "LL",
"value": 12.5
"value": 12.5,
"unit": "立方米/小时"
},
{
"key": "YL",
"value": 0.35
"value": 0.35,
"unit": "MPa"
}
]
}
```
### 水质数据
### 水质数据格式
```json
{
"testType": "常规检测",
@@ -145,68 +167,104 @@ tda:
}
```
## 监控和日志
## API 接口
### 日志级别
- **DEBUG**: 详细的数据处理日志
- **INFO**: 关键操作和状态变更
- **WARN**: 异常但可恢复的情况
- **ERROR**: 严重错误
### 数据采集管理
- `POST /api/data/collect/realtime` - 实时数据接入
- `POST /api/data/collect/batch` - 批量数据采集
- `GET /api/data/tasks` - 查询采集任务列表
- `GET /api/data/records` - 查询采集记录
### 监控指标
- **采集成功率**: 成功处理的数据量 / 总数据量
- **数据处理延迟**: 从接收到存储的时间差
- **系统负载**: CPU、内存使用率
- **连接状态**: MQTT、Kafka、数据库连接数
### 统计分析
- `GET /api/data/statistics` - 数据统计
- `GET /api/data/devices/{deviceSn}/statistics` - 设备统计
- `GET /api/data/errors/statistics` - 错误统计
## 开发指南
### MQTT 控制接口
- `POST /api/mqtt/control` - 发送设备控制命令
- `POST /api/mqtt/config` - 更新设备配置
### 添加新的数据源
1. 在 `DataCollectService` 中添加新的处理方法
2. 在 `DataValidationUtils` 中添加验证规则
3. 更新 `MetricType` 枚举(如需要)
4. 添加对应的单元测试
## WebSocket 主题
### 添加新的数据类型
1. 定义新的消息格式
2. 更新 Kafka 消费者
3. 添加数据验证逻辑
4. 实现存储逻辑
### 实时数据推送
- `/topic/data/realtime` - 全量实时数据
- `/topic/data/realtime/iot` - IoT 设备数据
- `/topic/data/realtime/quality` - 水质数据
### 控制指令
- `/topic/data/control` - 控制状态反馈
### 告警信息
- `/topic/data/alert` - 数据告警推送
### 统计数据
- `/topic/data/statistics` - 统计数据推送
## 测试
### 运行单元测试
### 单元测试
- `DataCollectServiceTest` - 数据采集服务测试
- `KafkaConsumerTest` - Kafka 消费者测试
- `DataValidationUtilsTest` - 数据验证测试
### 集成测试
- 实际 Kafka 服务器测试
- 实际 MQTT Broker 测试
- 数据库集成测试
## 部署和使用
### 环境要求
- Java 17+
- Spring Boot 3.3.5
- PostgreSQL 14+
- TDengine 3.0+
- Kafka 3.x+
- MQTT Broker (Eclipse Paho)
### 启动服务
```bash
mvn test
mvn spring-boot:run
```
### 运行集成测试
```bash
mvn verify
```
### 监控指标
- 数据采集成功率
- 数据处理延迟
- 内存使用情况
- 线程池状态
## 故障排查
## 问题排查
### 常见问题
1. **MQTT 连接失败**: 检查 Broker 地址和认证信息
2. **Kafka 消费延迟**: 检查消费者组和 Topic 配置
3. **TDengine 写入失败**: 检查数据库连接和超级表结构
4. **数据验证失败**: 检查数据格式和范围规则
1. **Kafka 连接失败**: 检查 Kafka 服务器地址和端口
2. **MQTT 连接失败**: 检查 Broker URL、用户名和密码
3. **TDengine 写入失败**: 检查数据库连接和表结构
4. **数据验证失败**: 检查数据格式和数值范围
### 调试模式
设置日志级别为 DEBUG:
### 日志配置
```yaml
logging:
level:
com.water.data_engine: DEBUG
org.springframework.kafka: INFO
org.eclipse.paho.client.mqttv3: WARN
```
## 版本历史
## 开发指南
### v1.0.0 (2026-06-15)
- 实现 Issue #41: 实时流数据采集(MQTT/Kafka Consumer)
- 支持 MQTT 客户端和数据接收
- 实现 Kafka 消费者功能
- 添加数据验证和质量检查
- 集成 TDengine 时序数据库
- 完善测试覆盖
### 添加新的数据源类型
1. 在 `MetricType` 枚举中添加新的指标类型
2. 在 `DataValidationUtils` 中添加对应的验证规则
3. 在 `DataCollectService` 中添加对应的处理逻辑
4. 更新配置文件中的 topic 路由规则
### 扩展数据验证规则
1. 在 `DataValidationUtils` 中添加新的验证方法
2. 在 `validateData()` 方法中调用新的验证逻辑
3. 编写对应的单元测试
### 添加新的数据存储后端
1. 实现新的存储接口
2. 在 `DataCollectService` 中集成新的存储后端
3. 添加配置选项
4. 编写集成测试