场景示例:MQTT 多点拆分写入 IoTDB
工业网关常把多个测点打在一条 MQTT JSON 里上报。业务侧需要:
- 拆成单测点消息,便于按测点订阅
- 用映射表把设备测点名换成业务测点名
- 转发给 MQTT 客户端,同时写入 IoTDB
控制台一条管道只能配置一个 Sink,因此拆成两条管道:管道 A 拆分后发到 MQTT;管道 B 写入 IoTDB。两条管道之间不直接用 MQTT 对接,而是用数据流把拆分后的 MQTT 主题沉淀为 Kafka Topic,管道 B 以 Kafka 为 Source。
链路
网关 --多测点报文--> 管道 A(MQTT Source:gateway/site1/#)
拆分数组 + MySQL 测点映射
MQTT Sink:unit1/{source}
│
┌───────────────┼───────────────┐
▼ ▼
MQTT 订阅者(unit1/#) 数据流 unit1-telemetry
按测点实时消费 主题过滤器:unit1/#
│
▼
管道 B(Kafka Source:unit1-telemetry)
IoTDB Sink:root.plant.unit1- 管道 A 只订网关主题
gateway/site1/#,不订unit1/#,以免把拆分结果再次拆分 - 数据流用主题过滤器捕获
unit1/#,无需再搭一套 Kafka - 管道 B 消费该数据流对应的 Kafka Topic,Processors 可留空,Sink 仍为 IoTDB
报文形态
网关报文(拆分前)
主题:gateway/site1/{deviceId},QoS 1。
data 为数组:每个元素带一个 timestamp 和若干测点字段。下面示例为 2 组 × 5 测点;完整报文一般为 10 组 × 5 测点 = 50 个数据点,拆分后对应 50 条单测点消息。
{
"id": "e6eb1d31-5ab7-4ace-8bf7-e7c3e7a30a73",
"source": "",
"mode": "5",
"type": "1",
"time": "2026-03-24T16:31:12.000+08:00",
"data": [
{
"timestamp": 1774341041606,
"DEV001.PT01": 281.1626812,
"DEV001.PT02": 281.2626822,
"DEV001.PT03": 281.3626832,
"DEV001.PT04": 281.4626842,
"DEV001.PT05": 281.5626852
},
{
"timestamp": 1774341041706,
"DEV001.PT06": 281.1626812,
"DEV001.PT07": 281.2626822,
"DEV001.PT08": 281.3626832,
"DEV001.PT09": 281.4626842,
"DEV001.PT10": 281.5626852
}
]
}单测点报文(拆分后)
主题:unit1/{source},例如 unit1/UNIT1.P01。
{
"id": "e6eb1d31-5ab7-4ace-8bf7-e7c3e7a30a73",
"source": "UNIT1.P01",
"time": "2026-03-24T16:31:12.000+08:00",
"type": "1",
"data": {
"value": 281.1626812,
"timestamp": 1774341041606
}
}| 字段 | 含义 |
|---|---|
id | 网关报文 ID,拆分后保持不变 |
source | 映射后的业务测点名,同时作为 MQTT 主题后缀和 IoTDB measurement |
time | 网关报文时间 |
type | 类型代码 |
data.value | 测点数值 |
data.timestamp | 时序时间戳(Unix 毫秒) |
测点映射表
在 MySQL 中维护设备测点到业务测点的对照。查不到映射的测点应丢弃,避免把设备侧原名发到业务主题。
CREATE DATABASE IF NOT EXISTS `point_map`;
CREATE TABLE `point_map`.`point_mapping` (
`id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
`device_point_field` VARCHAR(64) NOT NULL,
`biz_point_field` VARCHAR(64) NOT NULL,
PRIMARY KEY (`device_point_field`),
UNIQUE KEY `uk_point_mapping_id` (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
INSERT INTO `point_map`.`point_mapping` (`device_point_field`, `biz_point_field`) VALUES
('DEV001.PT01', 'UNIT1.P01'),
('DEV001.PT02', 'UNIT1.P02'),
('DEV001.PT03', 'UNIT1.P03'),
('DEV001.PT04', 'UNIT1.P04'),
('DEV001.PT05', 'UNIT1.P05'),
('DEV001.PT06', 'UNIT1.P06'),
('DEV001.PT07', 'UNIT1.P07'),
('DEV001.PT08', 'UNIT1.P08'),
('DEV001.PT09', 'UNIT1.P09'),
('DEV001.PT10', 'UNIT1.P10');管道 A:拆分并转发 MQTT
在控制台创建管道,名称例如 gateway-split-mqtt。
Source(MQTT)
| 参数 | 示例 |
|---|---|
| urls | tcp://127.0.0.1:1883 |
| topics | $share/split/gateway/site1/# |
| qos | 1 |
| clean_session | true |
多实例时使用共享订阅,避免重复处理同一条网关报文。client_id 需保证唯一(可开启动态后缀)。
Processors
把 data 数组拆成单测点,再 sql_select 查映射表,用 branch 把查询结果写回 source(直接 sql_select 会覆盖整条消息)。time / type 是表达式保留字,不能在 map_each 里用 let msg_time = this.time(会变成 null),需先写到非保留字段再还原。
processors:
- mapping: |
root.msg_id = this.id
root.msg_time = this.time
root.msg_type = this.type
root.points = this.data.map_each(item -> item.without("timestamp").key_values().map_each(kv -> {
"source": kv.key,
"data": {
"value": kv.value,
"timestamp": item.timestamp
}
})).flatten()
- mapping: |
root = this.points.map_each(p -> p.merge({
"saved_id": json("msg_id"),
"saved_time": json("msg_time"),
"saved_type": json("msg_type")
}))
- unarchive:
format: json_array
- mapping: |
root = this
root.id = this.saved_id
root.time = this.saved_time
root.type = this.saved_type
root.saved_id = deleted()
root.saved_time = deleted()
root.saved_type = deleted()
- branch:
processors:
- cached:
key: ${! json("source") }
cache: point_mapping_cache
ttl: 1h
skip_on: errored()
processors:
- sql_select:
driver: mysql
dsn: user:pass@tcp(127.0.0.1:3306)/point_map?parseTime=true
table: point_mapping
columns: [ biz_point_field ]
where: device_point_field = ?
args_mapping: root = [ this.source ]
result_map: |
root.source = this.index(0).biz_point_field
- catch:
- mapping: root = deleted()并在管道完整 YAML 顶层增加内存缓存(与 input / pipeline / output 同级,不要放进控制台 Processors 的 processors 列表):
cache_resources:
- label: point_mapping_cache
memory:
default_ttl: 1h表达式语法见 Processors。
Sink(MQTT)
| 参数 | 示例 |
|---|---|
| urls | tcp://127.0.0.1:1883 |
| topic | unit1/${! json("source") } |
| qos | 1 |
| retained | false |
映射后 source 为 UNIT1.P01 时,实际发布主题为 unit1/UNIT1.P01。MQTT 客户端可直接订 unit1/# 或按测点订 unit1/UNIT1.P01。
数据流:MQTT 拆分结果落到 Kafka
管道 A 发出的是 MQTT 主题。在控制台 数据流 中创建 Stream(同时作为 Kafka Topic),并绑定主题过滤器,把匹配的 MQTT 消息持久化进该流。之后 Kafka 客户端和管道 B 都按 Topic 消费,不必再订 MQTT。
- 打开「数据流」,创建 Stream,名称例如
unit1-telemetry(即后续 Kafka Topic 名) - 为该 Stream 添加主题过滤器:
unit1/# - 保存。此后发到
unit1/UNIT1.P01等主题的消息都会写入unit1-telemetry
主题过滤器与 MQTT 通配规则一致(+、#)。说明见 数据流 · 主题过滤器绑定。
管道 B:从数据流写入 IoTDB
再创建一条管道,名称例如 telemetry-to-iotdb。Processors 可留空。
Source(Kafka)
消费上一步创建的数据流。seed_brokers 填 FlowMQ 的 Kafka 入口。
| 参数 | 示例 |
|---|---|
| seed_brokers | 127.0.0.1:9092 |
| topics | unit1-telemetry |
| consumer_group | telemetry-to-iotdb |
多任务实例共用同一个 consumer_group,由组内协调分区。Kafka Source 参数见 Kafka Source。
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,以便小批量也能及时刷盘。
本场景测点名含 .(如 UNIT1.P01),写入 IoTDB 时必须给 measurement 加上反引号,否则会报 ILLEGAL_PATH。路径节点同样不能为纯数字;必须使用时写成 root.plant.`1`。
更多参数见 IoTDB Sink。
验证
- MQTT 客户端订
unit1/#(或$share/group/unit1/UNIT1.P01等按测点共享订阅) - 向
gateway/site1/device1发布一条完整网关报文(50 个数据点) - 订阅端应收到 50 条单测点消息,
source为UNIT1.P01…UNIT1.P10,且含id/time/type/data - 数据流
unit1-telemetry中应有这 50 条消息(Kafka 消费同一 Topic 可见) - 在 IoTDB 中确认:
SHOW TIMESERIES root.plant.unit1.**;
SELECT * FROM root.plant.unit1;预期 10 条时序(每个业务测点一条),完整 50 点报文对应每个测点 5 个时间戳。
注意事项
- 管道 A 的 Source 不要订
unit1/#,否则会把拆分结果再次拆分 - 管道 B 用 Kafka Source 读数据流,不要再用 MQTT 订
unit1/#。MQTT 订阅者与主题过滤器可以同时接收unit1/#,互不影响 - 先创建数据流并绑定
unit1/#,再启动管道 A / B,否则拆分结果不会进入unit1-telemetry - MySQL 在控制台只能作为 Sink 写入;本场景的查表必须写在 Processors 的
sql_select中 sql_select请放在branch里,并用catch丢弃无映射的测点- 管道 A 多任务实例时,MQTT Source 使用共享订阅,且
client_id唯一 - 管道 B 多任务实例时共用同一个 Kafka
consumer_group