Skip to content

DolphinDB Sink 集成

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

前提条件

  • 已了解管道结构与创建流程,参考 数据集成;Sink 连接器总览见 Sink
  • DolphinDB 服务可访问
  • 已具备目标库表写入权限

配置步骤

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

配置示例

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

可参考以下参数示例:

参数建议值
urls127.0.0.1:8848
usernameadmin
password123456
directorydfs://flowmq
tablemqtt_kafka_events
memory_cache_size10000
subscribe.batch_size1
subscribe.batch_interval1(单位:秒)

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

namedata_typeis_partitionis_sortvaluevalue_source
mpCodeSTRINGfalsetruesourcejson path
valueTimeTIMESTAMPtruefalsedata.timestampjson path
valueFLOATfalsefalsedata.valuejson path
unitCodeSTRINGfalsefalsetypejson path

每列指定 value_source(默认 direct)。上表示例为 json path。分区列 valueTime 映射 data.timestamp(毫秒),请确保落在库的日期分区范围内。样例 data.value 为字符串时也可写入 FLOAT 列。

value_source 说明(DolphinDB)

columns 中仅 valuevalue_source 控制;namedata_typeis_partitionis_sort 不受影响。

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

使用 Kafka Source(kafka_franz)时,可用的 metadata 键包括:kafka_keykafka_topickafka_partitionkafka_offsetkafka_timestamp_unixkafka_tombstone_message

必要表单参数

参数名后带 * 表示控制台必填。「支持表达式」为「是」时,可使用 ${! ... } 按消息动态取值,详见 Processors

参数支持表达式说明
urls*DolphinDB 地址列表(如 127.0.0.1:8848
username*用户名
password*密码
directory*目标库路径(如 dfs://flowmq
table*目标表名
memory_cache_size*流表缓存大小
subscribe*batch_size / batch_interval(秒)
columns*列映射(含 value_source
max_in_flight*并发写入上限

参数建议

  • subscribe.batch_interval 单位为秒;需要尽快看到写入结果时,可将 batch_size 设为 1
  • columns 中每列指定 value_source(默认 direct);从样例 JSON 取字段可优先用 json path
  • 分区字段取值须落在库自动生成的日期分区范围内

注意事项

  • 完成连通性测试后,建议再写入样例消息确认数据已落入目标表
  • 请确认 DolphinDB 版本与 FlowMQ 内置客户端兼容