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:
+184
-126
@@ -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. 编写集成测试
|
||||
Reference in New Issue
Block a user