Skip to content

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

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

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

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

控制台一条管道只能配置一个 Sink,因此拆成两条管道:管道 A 拆分数组后用 Kafka Sink 写到 iot-hub-out;管道 B 从该 Topic 写入 IoTDB。Kafka 业务消费者与管道 B 订同一个 iot-hub-out,用不同消费组互不影响。

链路

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

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 和 Kafka key
TM时序时间戳(Unix 毫秒)
AV测点数值(字符串或数字均可,写入 IoTDB 前转为 FLOAT)

拆分后

管道 A 写入 Kafka Topic iot-hub-out,消息 key 为 source(例如 UNIT1.P000)。

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

数据流

在控制台 数据流 中创建两条 Stream(无需 MQTT 主题过滤器):

  1. iot-hub-in:给 Kafka 生产者写入原始数组
  2. iot-hub-out:给管道 A 写入拆分后的单测点消息

Topic 不存在时,连通性测试可能通过,但无法完成实际转发。详见 Kafka Sink

管道 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(Kafka Franz)

参数示例
seed_brokers127.0.0.1:9092
topiciot-hub-out
key${! json("source") }
client_idiot-hub-split-sink

sourceUNIT1.P000 时,Kafka key 为 UNIT1.P000。不要把 topic 设成 iot-hub-in,以免形成消息环路。更多参数见 Kafka Sink

管道 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 条,source / key 为 UNIT1.P000UNIT1.P099
  4. 在 IoTDB 中确认:
sql
SHOW TIMESERIES root.plant.unit1.*;
SELECT COUNT(*) FROM root.plant.unit1;

预期 100 条时序,每个测点 30 个时间戳。

注意事项

  • 管道 A 只消费 iot-hub-in,Sink 只写 iot-hub-out,以免把拆分结果再次拆分
  • 先创建 iot-hub-iniot-hub-out,再启动管道 A / B
  • 本场景没有 MySQL 测点映射;若 CN 仍是设备侧原名,做法见 MQTT 多点拆分
  • 管道 A / B 多任务实例时,各自共用同一个 Kafka consumer_group