DolphinDB Sink 集成
本文介绍如何将 FlowMQ 数据写入 DolphinDB。
前提条件
配置步骤
- 在“数据管道”中创建管道,完成基础信息配置。
- 在 Source 步骤选择
Kafka(Topic:flowmq.mqtt.kafka),Processors 按需配置。 - 在 Sink 步骤选择
DolphinDB,配置urls、directory、table、columns等参数。 - 点击“测试”验证连接。
- 继续完成确认步骤。
配置示例
以下示例将 Kafka Topic flowmq.mqtt.kafka 中的样例消息写入目标系统。样例 JSON 与 Source 约定见 Sink · 示例场景。
可参考以下参数示例:
| 参数 | 建议值 |
|---|---|
| urls | 127.0.0.1:8848 |
| username | admin |
| password | 123456 |
| directory | dfs://flowmq |
| table | mqtt_kafka_events |
| memory_cache_size | 10000 |
| subscribe.batch_size | 1 |
| subscribe.batch_interval | 1(单位:秒) |
columns 字段映射示例(对应样例消息):
| name | data_type | is_partition | is_sort | value | value_source |
|---|---|---|---|---|---|
| mpCode | STRING | false | true | source | json path |
| valueTime | TIMESTAMP | true | false | data.timestamp | json path |
| value | FLOAT | false | false | data.value | json path |
| unitCode | STRING | false | false | type | json path |
每列指定
value_source(默认direct)。上表示例为json path。分区列valueTime映射data.timestamp(毫秒),请确保落在库的日期分区范围内。样例data.value为字符串时也可写入FLOAT列。
value_source 说明(DolphinDB)
columns 中仅 value 受 value_source 控制;name、data_type、is_partition、is_sort 不受影响。
| value_source | value 填写内容 | 运行时行为 |
|---|---|---|
direct | 固定字面量,或 ${! ... } 表达式 | 不含 ${!} 作为静态值;含 ${! ... } 时按每条消息执行插值 |
meta | 元数据键名(字符串) | 从消息 metadata 取对应键,键不存在时报错 |
json path | JSON 点路径(如 source、data.value、data.sensors[0].temp) | 将消息解析为结构化 JSON 后按路径取值,路径不存在或值为 null 时报错([N] 解析为 .N) |
使用 Kafka Source(kafka_franz)时,可用的 metadata 键包括:kafka_key、kafka_topic、kafka_partition、kafka_offset、kafka_timestamp_unix、kafka_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设为1columns中每列指定value_source(默认direct);从样例 JSON 取字段可优先用json path- 分区字段取值须落在库自动生成的日期分区范围内
注意事项
- 完成连通性测试后,建议再写入样例消息确认数据已落入目标表
- 请确认 DolphinDB 版本与 FlowMQ 内置客户端兼容