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

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

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

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

与 [MQTT 多点拆分](data-integration-scenario-mqtt-split-iotdb.md)（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 消息

主题过滤器说明见 [数据流 · 主题过滤器绑定](streaming.md#主题过滤器绑定)。

## 管道 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`：

```yaml
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）

```sql
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](data-integration-sink-iotdb.md)。

## 验证

1. Kafka 消费者订 `iot-hub-out`（独立消费组）
2. 向 `iot-hub-in` 发布一条 3000 点数组
3. 消费者应收到 **3000** 条，`source` 为 `UNIT1.P000` … `UNIT1.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 多点拆分](data-integration-scenario-mqtt-split-iotdb.md)
- 管道 A / B 多任务实例时，各自共用同一个 Kafka `consumer_group`
