IoTDB Sink 集成
本文介绍如何将 FlowMQ 数据写入 Apache IoTDB。
前提条件
配置步骤
- 在“数据管道”中创建管道,完成基础信息配置。
- 在 Source 步骤选择
Kafka(Topic:flowmq.mqtt.kafka),Processors 按需配置。 - 在 Sink 步骤选择
IoTDB,配置urls、device_id、columns等参数。 - 点击“测试”验证连接。
- 继续完成确认步骤。
配置示例
以下示例将 Kafka Topic flowmq.mqtt.kafka 中的样例消息写入目标系统。样例 JSON 与 Source 约定见 Sink · 示例场景。
可参考以下参数示例:
| 参数 | 建议值 |
|---|---|
| urls | 127.0.0.1:6667 |
| username | root |
| password | root |
| device_id | root.flowmq.demo |
| timestamp_precision | ms |
columns 字段映射示例(对应样例消息):
| timestamp | measurement | data_type | value | value_source |
|---|---|---|---|---|
${! json("data.timestamp") } | ${! json("source") } | FLOAT | ${! json("data.value") } | direct |
每列指定
value_source(默认direct);同一列的timestamp/measurement/value共用该方式。使用direct时填写${! json("...") }从样例 JSON 取值。样例data.timestamp为毫秒时,请将timestamp_precision设为ms。
value_source 说明(IoTDB)
在 IoTDB 的 columns 中,同一列的 timestamp、measurement、value 共用一个 value_source。data_type 仍是独立类型枚举字段,不受 value_source 控制。
| value_source | timestamp / measurement / value 填写内容 | 运行时行为 |
|---|---|---|
direct | 固定字面量,或 ${! ... } 表达式 | 不含 ${!} 作为静态值;含 ${! ... } 时按每条消息执行插值 |
meta | 元数据键名(字符串) | 从消息 metadata 取对应键,键不存在时报错 |
json path | JSON 点路径(如 data.timestamp、source、data.value) | 将消息解析为结构化 JSON 后按路径取值,路径不存在或值为 null 时报错([N] 解析为 .N) |
device_id 是顶层字段,支持插值字符串,但不使用 value_source 枚举。
使用 Kafka Source(kafka_franz)时,可用的 metadata 键包括:kafka_key、kafka_topic、kafka_partition、kafka_offset、kafka_timestamp_unix、kafka_tombstone_message。
必要表单参数
参数名后带 * 表示控制台必填。「支持表达式」为「是」时,可使用 ${! ... } 按消息动态取值,详见 Processors。
| 参数 | 支持表达式 | 说明 |
|---|---|---|
urls* | 否 | IoTDB 地址列表(如 127.0.0.1:6667) |
username* | 否 | 用户名 |
password* | 否 | 密码 |
device_id* | 是 | 设备路径 |
columns* | 是 | 列映射(含 value_source) |
max_in_flight* | 否 | 并发写入上限 |
timestamp_precision | 否 | 时间精度(ms / us / ns) |
参数建议
- 统一规划
device_id命名层级,避免路径混乱 columns中每列指定value_source(默认direct);从样例 JSON 取字段时用${! json("...") }timestamp_precision与消息时间戳单位保持一致