TDengine Sink 集成
本文介绍如何将 FlowMQ 数据写入 TDengine。
前提条件
- 已了解管道结构与创建流程,参考 数据集成;Sink 连接器总览见 Sink
- TDengine 服务可访问(WebSocket 端口,默认
6041) - TDengine 服务端版本建议 ≥
3.3.6.0(当前驱动要求) - 已具备目标库表写入权限
配置步骤
- 在“数据管道”中创建管道,完成基础信息配置。
- 在 Source 步骤选择
Kafka(Topic:flowmq.mqtt.kafka),Processors 按需配置。 - 在 Sink 步骤选择
TDengine,配置host、database、table、columns等参数。 - 点击“测试”验证连接。
- 继续完成确认步骤。
配置示例
以下示例将 Kafka Topic flowmq.mqtt.kafka 中的样例消息写入目标系统。样例 JSON 与 Source 约定见 Sink · 示例场景。
可参考以下参数示例:
| 参数 | 建议值 |
|---|---|
| host | 127.0.0.1 |
| port | 6041 |
| username | root |
| password | taosdata(默认安装密码,以实际为准) |
| database | flowmq |
| table | mqtt_kafka_events |
目标表示例(普通表;首列须为 TIMESTAMP):
sql
CREATE DATABASE IF NOT EXISTS flowmq PRECISION 'ms';
USE flowmq;
CREATE TABLE IF NOT EXISTS mqtt_kafka_events (
create_time TIMESTAMP,
source NCHAR(128),
value_time BIGINT,
`value` FLOAT,
unit_code NCHAR(64)
);columns 字段映射示例(对应样例消息):
| name | data_type | is_primary_key | value |
|---|---|---|---|
| create_time | TIMESTAMP | false | ${! json("time").ts_parse("2006-01-02T15:04:05.999Z07:00") } |
| source | VARCHAR | false | ${! json("source") } |
| value_time | BIGINT | false | ${! json("data.timestamp") } |
| value | FLOAT | false | ${! json("data.value") } |
| unit_code | VARCHAR | false | ${! json("type") } |
data.timestamp为 Unix 毫秒,建议以BIGINT落库。写入TIMESTAMP列时请用ts_parse(或${! now() }),勿直接填入 ISO 字符串。value为保留字,建表需加反引号。普通表首列TIMESTAMP具有唯一性,重复时间戳会导致写入失败。
必要表单参数
参数名后带 * 表示控制台必填。「支持表达式」为「是」时,可使用 ${! ... } 按消息动态取值,详见 Processors。
| 参数 | 支持表达式 | 说明 |
|---|---|---|
host* | 否 | 数据库地址 |
port* | 否 | WebSocket 端口(默认 6041,非原生 6030) |
username* | 否 | 数据库用户名 |
password* | 否 | 数据库密码 |
database* | 否 | 目标数据库 |
table* | 否 | 目标表 |
columns* | 是 | 列映射(name / data_type / is_primary_key / value) |
max_in_flight* | 否 | 并发写入上限 |
参数建议
- 预先建好目标表;首列使用
TIMESTAMP,精度与库配置一致(示例为毫秒) - 标签/列名避免使用保留字;需要时用反引号或改名
注意事项
- 端口请填 WebSocket 端口(常见为
6041),原生端口6030不可用于本 Sink - 服务端版本过低时,运行期可能报版本不匹配(要求 ≥
3.3.6.0) - 毫秒时间戳字段不要直接映射为
TIMESTAMP;ISO 字符串请用ts_parse后再写入TIMESTAMP列