场景示例:Kafka 多点拆分写入 IoTDB
物联中心常把同一批次的大量测点打成一条 Kafka JSON 数组上报。业务侧需要:
- 拆成单测点消息
- 转发给 Kafka 消费者
- 同时写入 IoTDB
与 MQTT 多点拆分(MQTT 网关入口 + MySQL 测点改名 + MQTT Sink)不同:本场景入口和拆分后的转发都走 Kafka,测点名已在报文的 CN 中,无需再查映射表。
控制台一条管道只能配置一个 Sink,因此拆成两条管道:管道 A 拆分数组后用 Kafka Sink 写到 iot-hub-out;管道 B 从该 Topic 写入 IoTDB。Kafka 业务消费者与管道 B 订同一个 iot-hub-out,用不同消费组互不影响。
链路
Kafka 生产者 --多测点数组--> 数据流 iot-hub-in
│
管道 A(Kafka Source:iot-hub-in)
unarchive 数组 + CN/AV/TM 映射
Kafka Sink:iot-hub-out
│
数据流 iot-hub-out
│
┌────────────────────┼────────────────────┐
▼ ▼
Kafka 消费者(iot-hub-out) 管道 B(Kafka Source:iot-hub-out)
IoTDB Sink:root.plant.unit1- 生产者直接写数据流
iot-hub-in(Kafka Topic) - 管道 A 只消费
iot-hub-in,Sink 写iot-hub-out,不要写回iot-hub-in - Kafka 消费者与管道 B 都读
iot-hub-out
报文形态
拆分前(Payload3)
写入 Kafka Topic iot-hub-in。根是 JSON 数组,每个元素一个数据点:CN 测点名、TM 时间戳(Unix 毫秒)、AV 数值。完整报文一般为 100 个测点 × 30 个时间戳 = 3000 个数据点,拆分后对应 3000 条单测点消息。
[
{"CN": "UNIT1.P000", "TM": 1774341041606, "AV": 0},
{"CN": "UNIT1.P001", "TM": 1774341041606, "AV": 1},
{"CN": "UNIT1.P000", "TM": 1774341041706, "AV": 0.001}
]| 字段 | 含义 |
|---|---|
CN | 业务测点名,拆分后作为 source 和 Kafka key |
TM | 时序时间戳(Unix 毫秒) |
AV | 测点数值(字符串或数字均可,写入 IoTDB 前转为 FLOAT) |
拆分后
管道 A 写入 Kafka Topic iot-hub-out,消息 key 为 source(例如 UNIT1.P000)。
{
"source": "UNIT1.P000",
"data": {
"value": 0,
"timestamp": 1774341041606
}
}数据流
在控制台 数据流 中创建两条 Stream(无需 MQTT 主题过滤器):
iot-hub-in:给 Kafka 生产者写入原始数组iot-hub-out:给管道 A 写入拆分后的单测点消息
Topic 不存在时,连通性测试可能通过,但无法完成实际转发。详见 Kafka Sink。
管道 A:拆分并转发
名称例如 iot-hub-split。
Source(Kafka)
| 参数 | 示例 |
|---|---|
| seed_brokers | 127.0.0.1:9092 |
| topics | iot-hub-in |
| consumer_group | iot-hub-split |
多任务实例共用同一个 consumer_group。
Processors
根是数组,先 unarchive 成单条,再把 CN/AV/TM 改成下游使用的 source + data:
processors:
- unarchive:
format: json_array
- mapping: |
root.source = this.CN
root.data = {
"value": this.AV.number(),
"timestamp": this.TM
}Sink(Kafka Franz)
| 参数 | 示例 |
|---|---|
| seed_brokers | 127.0.0.1:9092 |
| topic | iot-hub-out |
| key | ${! json("source") } |
| client_id | iot-hub-split-sink |
source 为 UNIT1.P000 时,Kafka key 为 UNIT1.P000。不要把 topic 设成 iot-hub-in,以免形成消息环路。更多参数见 Kafka Sink。
管道 B:从数据流写入 IoTDB
名称例如 iot-hub-to-iotdb。Processors 可留空。
Source(Kafka)
| 参数 | 示例 |
|---|---|
| seed_brokers | 127.0.0.1:9092 |
| topics | iot-hub-out |
| consumer_group | iot-hub-to-iotdb |
Kafka 业务消费者使用另一个 consumer_group 订同一个 iot-hub-out,与管道 B 互不影响。
Sink(IoTDB)
CREATE DATABASE root.plant;| 参数 | 示例 |
|---|---|
| urls | 127.0.0.1:6667 |
| username / password | root / root |
| device_id | root.plant.unit1 |
| timestamp_precision | ms |
columns 映射:
| timestamp | measurement | data_type | value | value_source |
|---|---|---|---|---|
${! json("data.timestamp") } | `${! json("source") }` | FLOAT | ${! json("data.value") } | direct |
batching.byte_size 须为整数(例如 4194304),不要写成 4MB。建议同时设置 count 与 period。
测点名含 . 时,measurement 需用反引号包裹,否则 IoTDB 报 ILLEGAL_PATH。路径节点不能为纯数字;必须使用时写成 root.plant.`1`。更多参数见 IoTDB Sink。
验证
- Kafka 消费者订
iot-hub-out(独立消费组) - 向
iot-hub-in发布一条 3000 点数组 - 消费者应收到 3000 条,
source/ key 为UNIT1.P000…UNIT1.P099 - 在 IoTDB 中确认:
SHOW TIMESERIES root.plant.unit1.*;
SELECT COUNT(*) FROM root.plant.unit1;预期 100 条时序,每个测点 30 个时间戳。
注意事项
- 管道 A 只消费
iot-hub-in,Sink 只写iot-hub-out,以免把拆分结果再次拆分 - 先创建
iot-hub-in与iot-hub-out,再启动管道 A / B - 本场景没有 MySQL 测点映射;若
CN仍是设备侧原名,做法见 MQTT 多点拆分 - 管道 A / B 多任务实例时,各自共用同一个 Kafka
consumer_group