Skip to content

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

物联中心常把同一批次的大量测点打成一条 Kafka JSON 数组上报。业务侧需要:

  1. 拆成单测点消息
  2. 转发给 Kafka 消费者
  3. 同时写入 IoTDB

MQTT 多点拆分(MQTT 网关入口 + MySQL 测点改名)不同:本场景入口是 Kafka,测点名已在报文的 CN 中,无需再查映射表。

控制台一条管道只能配置一个 Sink,因此拆成两条管道:管道 A 拆分数组并发到 MQTT;数据流用主题过滤器把这些 MQTT 消息沉淀为 Kafka Topic,供 Kafka 消费者和管道 B 使用。

链路

text
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 条单测点消息。

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

json
{
  "source": "UNIT1.P000",
  "data": {
    "value": 0,
    "timestamp": 1774341041606
  }
}

数据流

在控制台 数据流 中创建两条 Stream:

  1. iot-hub-in:给 Kafka 生产者写入原始数组(无需额外 MQTT 过滤器)
  2. iot-hub-out:添加主题过滤器 telemetry/#,捕获管道 A 拆分后发布的 MQTT 消息

主题过滤器说明见 数据流 · 主题过滤器绑定

管道 A:拆分并转发

名称例如 iot-hub-split

Source(Kafka)

参数示例
seed_brokers127.0.0.1:9092
topicsiot-hub-in
consumer_groupiot-hub-split

多任务实例共用同一个 consumer_group

Processors

根是数组,先 unarchive 成单条,再把 CN/AV/TM 改成下游使用的 source + data

yaml
processors:
  - unarchive:
      format: json_array
  - mapping: |
      root.source = this.CN
      root.data = {
        "value": this.AV.number(),
        "timestamp": this.TM
      }

Sink(MQTT)

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

sourceUNIT1.P000 时,发布主题为 telemetry/UNIT1.P000

管道 B:从数据流写入 IoTDB

名称例如 iot-hub-to-iotdb。Processors 可留空。

Source(Kafka)

参数示例
seed_brokers127.0.0.1:9092
topicsiot-hub-out
consumer_groupiot-hub-to-iotdb

Kafka 业务消费者使用另一个 consumer_group 订同一个 iot-hub-out,与管道 B 互不影响。

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

测点名含 . 时,measurement 需用反引号包裹,否则 IoTDB 报 ILLEGAL_PATH。路径节点不能为纯数字;必须使用时写成 root.plant.`1`。更多参数见 IoTDB Sink

验证

  1. Kafka 消费者订 iot-hub-out(独立消费组)
  2. iot-hub-in 发布一条 3000 点数组
  3. 消费者应收到 3000 条,sourceUNIT1.P000UNIT1.P099
  4. (可选)MQTT 订 telemetry/# 也应看到这 3000 条
  5. 在 IoTDB 中确认:
sql
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