Skip to content

场景示例:MQTT 多点拆分写入 IoTDB

工业网关常把多个测点打在一条 MQTT JSON 里上报。业务侧需要:

  1. 拆成单测点消息,便于按测点订阅
  2. 用映射表把设备测点名换成业务测点名
  3. 转发给 MQTT 客户端,同时写入 IoTDB

控制台一条管道只能配置一个 Sink,因此拆成两条管道:管道 A 拆分后发到 MQTT;管道 B 写入 IoTDB。两条管道之间不直接用 MQTT 对接,而是用数据流把拆分后的 MQTT 主题沉淀为 Kafka Topic,管道 B 以 Kafka 为 Source。

链路

text
网关  --多测点报文-->  管道 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 条单测点消息。

json
{
  "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

json
{
  "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 中维护设备测点到业务测点的对照。查不到映射的测点应丢弃,避免把设备侧原名发到业务主题。

sql
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)

参数示例
urlstcp://127.0.0.1:1883
topics$share/split/gateway/site1/#
qos1
clean_sessiontrue

多实例时使用共享订阅,避免重复处理同一条网关报文。client_id 需保证唯一(可开启动态后缀)。

Processors

data 数组拆成单测点,再 sql_select 查映射表,用 branch 把查询结果写回 source(直接 sql_select 会覆盖整条消息)。time / type 是表达式保留字,不能在 map_each 里用 let msg_time = this.time(会变成 null),需先写到非保留字段再还原。

yaml
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 列表):

yaml
cache_resources:
  - label: point_mapping_cache
    memory:
      default_ttl: 1h

表达式语法见 Processors

Sink(MQTT)

参数示例
urlstcp://127.0.0.1:1883
topicunit1/${! json("source") }
qos1
retainedfalse

映射后 sourceUNIT1.P01 时,实际发布主题为 unit1/UNIT1.P01。MQTT 客户端可直接订 unit1/# 或按测点订 unit1/UNIT1.P01

数据流:MQTT 拆分结果落到 Kafka

管道 A 发出的是 MQTT 主题。在控制台 数据流 中创建 Stream(同时作为 Kafka Topic),并绑定主题过滤器,把匹配的 MQTT 消息持久化进该流。之后 Kafka 客户端和管道 B 都按 Topic 消费,不必再订 MQTT。

  1. 打开「数据流」,创建 Stream,名称例如 unit1-telemetry(即后续 Kafka Topic 名)
  2. 为该 Stream 添加主题过滤器:unit1/#
  3. 保存。此后发到 unit1/UNIT1.P01 等主题的消息都会写入 unit1-telemetry

主题过滤器与 MQTT 通配规则一致(+#)。说明见 数据流 · 主题过滤器绑定

管道 B:从数据流写入 IoTDB

再创建一条管道,名称例如 telemetry-to-iotdb。Processors 可留空。

Source(Kafka)

消费上一步创建的数据流。seed_brokers 填 FlowMQ 的 Kafka 入口。

参数示例
seed_brokers127.0.0.1:9092
topicsunit1-telemetry
consumer_grouptelemetry-to-iotdb

多任务实例共用同一个 consumer_group,由组内协调分区。Kafka Source 参数见 Kafka Source

Sink(IoTDB)

写入前先创建数据库:

sql
CREATE DATABASE root.plant;
参数示例
urls127.0.0.1:6667
username / passwordroot / root
device_idroot.plant.unit1
timestamp_precisionms

columns 映射:

timestampmeasurementdata_typevaluevalue_source
${! json("data.timestamp") }`${! json("source") }`FLOAT${! json("data.value") }direct

batching.byte_size 须为整数(例如 4194304),不要写成 4MB。建议同时设置 countperiod,以便小批量也能及时刷盘。

本场景测点名含 .(如 UNIT1.P01),写入 IoTDB 时必须给 measurement 加上反引号,否则会报 ILLEGAL_PATH。路径节点同样不能为纯数字;必须使用时写成 root.plant.`1`

更多参数见 IoTDB Sink

验证

  1. MQTT 客户端订 unit1/#(或 $share/group/unit1/UNIT1.P01 等按测点共享订阅)
  2. gateway/site1/device1 发布一条完整网关报文(50 个数据点)
  3. 订阅端应收到 50 条单测点消息,sourceUNIT1.P01UNIT1.P10,且含 id / time / type / data
  4. 数据流 unit1-telemetry 中应有这 50 条消息(Kafka 消费同一 Topic 可见)
  5. 在 IoTDB 中确认:
sql
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