Skip to content

TDengine Sink 集成

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

前提条件

  • 已了解管道结构与创建流程,参考 数据集成;Sink 连接器总览见 Sink
  • TDengine 服务可访问(WebSocket 端口,默认 6041
  • TDengine 服务端版本建议 ≥ 3.3.6.0(当前驱动要求)
  • 已具备目标库表写入权限

配置步骤

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

配置示例

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

可参考以下参数示例:

参数建议值
host127.0.0.1
port6041
usernameroot
passwordtaosdata(默认安装密码,以实际为准)
databaseflowmq
tablemqtt_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 字段映射示例(对应样例消息):

namedata_typeis_primary_keyvalue
create_timeTIMESTAMPfalse${! json("time").ts_parse("2006-01-02T15:04:05.999Z07:00") }
sourceVARCHARfalse${! json("source") }
value_timeBIGINTfalse${! json("data.timestamp") }
valueFLOATfalse${! json("data.value") }
unit_codeVARCHARfalse${! 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