场景示例:Kafka 多点拆分写入 IoTDB
物联中心常把同一批次的大量测点打成一条 Kafka JSON 数组上报。业务侧需要:
- 拆成单测点消息
- 转发给 Kafka 消费者
- 同时写入 IoTDB
与 MQTT 多点拆分(MQTT 网关入口 + MySQL 测点改名)不同:本场景入口是 Kafka,测点名已在报文的 CN 中,无需再查映射表。
控制台一条管道只能配置一个 Sink,因此拆成两条管道:管道 A 拆分数组并发到 MQTT;数据流用主题过滤器把这些 MQTT 消息沉淀为 Kafka Topic,供 Kafka 消费者和管道 B 使用。
链路
Kafka 生产者 --多测点数组--> 数据流 iot-hub-in
│
管道 A(Kafka Source:iot-hub-in)
unarchive 数组 + CN/AV/TM 映射
MQTT Sink:telemetry/{source}
│
┌─────────────┼─────────────┐
▼ ▼
(可选)MQTT 订 telemetry/# 数据流 iot-hub-out
主题过滤器:telemetry/#
│
┌────────────────────┼────────────────────┐
▼ ▼
Kafka 消费者(iot-hub-out) 管道 B(Kafka Source:iot-hub-out)
IoTDB Sink:root.plant.unit1- 生产者直接写数据流
iot-hub-in(Kafka Topic),不必再经 MQTT - 管道 A 只消费
iot-hub-in,不要订telemetry/# - Kafka 消费者与管道 B 都读
iot-hub-out,不要再用 MQTT Source 订telemetry/#
报文形态
拆分前(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 和 MQTT 主题后缀 |
TM | 时序时间戳(Unix 毫秒) |
AV | 测点数值(字符串或数字均可,写入 IoTDB 前转为 FLOAT) |
拆分后
MQTT 主题:telemetry/{source},例如 telemetry/UNIT1.P000。同一条消息也会进入数据流 iot-hub-out。
{
"source": "UNIT1.P000",
"data": {
"value": 0,
"timestamp": 1774341041606
}
}数据流
在控制台 数据流 中创建两条 Stream:
iot-hub-in:给 Kafka 生产者写入原始数组(无需额外 MQTT 过滤器)iot-hub-out:添加主题过滤器telemetry/#,捕获管道 A 拆分后发布的 MQTT 消息
主题过滤器说明见 数据流 · 主题过滤器绑定。
管道 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(MQTT)
| 参数 | 示例 |
|---|---|
| urls | tcp://127.0.0.1:1883 |
| topic | telemetry/${! json("source") } |
| qos | 1 |
| retained | false |
source 为 UNIT1.P000 时,发布主题为 telemetry/UNIT1.P000。
管道 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为UNIT1.P000…UNIT1.P099 - (可选)MQTT 订
telemetry/#也应看到这 3000 条 - 在 IoTDB 中确认:
SHOW TIMESERIES root.plant.unit1.*;
SELECT COUNT(*) FROM root.plant.unit1;预期 100 条时序,每个测点 30 个时间戳。
注意事项
- 管道 A 只消费
iot-hub-in,不要订telemetry/#,以免把拆分结果再次拆分 - 先创建
iot-hub-out并绑定telemetry/#,再启动管道 A / B,否则拆分结果不会进入数据流 - 本场景没有 MySQL 测点映射;若
CN仍是设备侧原名,做法见 MQTT 多点拆分 - 管道 A / B 多任务实例时,各自共用同一个 Kafka
consumer_group