Skip to content

IoTDB Sink 集成

本文介绍如何将 FlowMQ 数据写入 Apache IoTDB。

前提条件

  • 已了解管道结构与创建流程,参考 数据集成;Sink 连接器总览见 Sink
  • IoTDB 服务可访问(RPC 端口默认 6667
  • 已具备写入目标路径的权限

配置步骤

  1. 在“数据管道”中创建管道,完成基础信息配置。
  2. 在 Source 步骤选择 Kafka(Topic:flowmq.mqtt.kafka),Processors 按需配置。
  3. 在 Sink 步骤选择 IoTDB,配置 urlsdevice_idcolumns 等参数。
  4. 点击“测试”验证连接。
  5. 继续完成确认步骤。

配置示例

以下示例将 Kafka Topic flowmq.mqtt.kafka 中的样例消息写入目标系统。样例 JSON 与 Source 约定见 Sink · 示例场景

可参考以下参数示例:

参数建议值
urls127.0.0.1:6667
usernameroot
passwordroot
device_idroot.flowmq.demo
timestamp_precisionms

columns 字段映射示例(对应样例消息):

timestampmeasurementdata_typevaluevalue_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 中,同一列的 timestampmeasurementvalue 共用一个 value_sourcedata_type 仍是独立类型枚举字段,不受 value_source 控制。

value_sourcetimestamp / measurement / value 填写内容运行时行为
direct固定字面量,或 ${! ... } 表达式不含 ${!} 作为静态值;含 ${! ... } 时按每条消息执行插值
meta元数据键名(字符串)从消息 metadata 取对应键,键不存在时报错
json pathJSON 点路径(如 data.timestampsourcedata.value将消息解析为结构化 JSON 后按路径取值,路径不存在或值为 null 时报错([N] 解析为 .N

device_id 是顶层字段,支持插值字符串,但不使用 value_source 枚举。

使用 Kafka Source(kafka_franz)时,可用的 metadata 键包括:kafka_keykafka_topickafka_partitionkafka_offsetkafka_timestamp_unixkafka_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 与消息时间戳单位保持一致