TDengine 与 Kafka 双向数据同步实战:Kafka Connect Source/Sink Connector 完全指南
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
TDengine Kafka Connector 是面向 Kafka Connect 生态的一对插件(TDengine Source Connector 与 TDengine Sink Connector),只需一份简单的 JSON 配置即可实现 Kafka 指定 topic 与 TDengine 指定数据库之间的批量或实时双向同步。本文将以官方文档为主线,结合仓库中的无模式(Schemaless)写入文档与数据订阅文档,完整演示从环境搭建、插件编译安装、Sink/Source 两个方向的端到端实战,并给出全部配置参数的参考说明,帮助你在真实 IIoT 场景中快速搭建 Kafka 与 TDengine 之间的数据管道。
什么是 Kafka Connect?
Kafka Connect 是 Apache Kafka 的一个组件,用于让其他系统(数据库、云服务、文件系统等)方便地接入 Kafka。数据既可以通过 Kafka Connect 从外部系统流入 Kafka,也可以从 Kafka 流向外部系统。其中:
- Source Connector:从其他系统读取数据,交给 Kafka Connect 写入 Kafka;
- Sink Connector:从 Kafka Connect 接收 Kafka 中的数据,写入其他系统。
需要注意的是,Source Connector 和 Sink Connector 都不会直接连接 Kafka Broker——Source Connector 把数据转交给 Kafka Connect,Sink Connector 从 Kafka Connect 接收数据。
TDengine Kafka Connector 基于这一模型实现了两个插件:
- TDengine Source Connector:实时从 TDengine 读取数据并发送给 Kafka Connect;
- TDengine Sink Connector:从 Kafka Connect 接收数据并写入 TDengine。
TDengine Source Connector 与 Sink Connector 共同构成了 TDengine 与 Kafka 之间的双向数据通道,可同时支撑"Kafka 数据汇入 TDengine"(如消息总线落库)与"TDengine 数据发布到 Kafka"(如实时数仓入湖)两类典型场景。
前置条件
运行本教程示例需要满足以下环境要求:
- Linux 操作系统;
- 已安装 Java 8 和 Maven;
- 已安装 Git、curl、vi;
- 已安装并启动 TDengine。如果尚未安装,可参考安装与卸载指南。
说明:示例中的连接串
jdbc:TAOS://127.0.0.1:6030使用 TDengine 默认服务端口 6030,请确保本机 TDengine 已正常运行且root/taosdata账号可登录。
安装 Kafka
在任意目录下执行以下命令,下载并解压 Kafka 3.4.0 发行包:
KAFKA_PKG="kafka_2.13-3.4.0" curl -O "https://archive.apache.org/dist/kafka/3.4.0/${KAFKA_PKG}.tgz" tar xzf "${KAFKA_PKG}.tgz" -C /opt/ ln -s "/opt/${KAFKA_PKG}" /opt/kafka随后将$KAFKA_HOME/bin目录加入 PATH。将以下内容追加到当前用户的 profile 文件(~/.profile或~/.bash_profile):
export KAFKA_HOME=/opt/kafka export PATH=$PATH:$KAFKA_HOME/bin保存后执行source ~/.profile(或重新登录终端)使环境变量生效。
安装 TDengine Connector 插件
编译插件
git clone --branch 3.0 https://github.com/taosdata/kafka-connect-tdengine.git cd kafka-connect-tdengine mvn clean package -Dmaven.test.skip=true unzip -d $KAFKA_HOME/components/ target/components/packages/taosdata-kafka-connect-tdengine-*.zip以上脚本首先克隆项目源码,然后使用 Maven 编译打包。打包完成后,插件的 zip 包生成在target/components/packages/目录中,将其解压到插件安装路径即可。示例中使用的是 Kafka 内置的插件安装路径$KAFKA_HOME/components/。
配置插件
编辑$KAFKA_HOME/config/connect-distributed.properties,将 kafka-connect-tdengine 插件目录加入plugin.path:
plugin.path=/usr/share/java,/opt/kafka/components启动 Kafka 与 Kafka Connect
依次启动 ZooKeeper、Kafka Broker 与 Kafka Connect(分布式模式):
zookeeper-server-start.sh -daemon $KAFKA_HOME/config/zookeeper.properties kafka-server-start.sh -daemon $KAFKA_HOME/config/server.properties connect-distributed.sh -daemon $KAFKA_HOME/config/connect-distributed.properties验证 Kafka Connect 是否启动成功
Kafka Connect 分布式模式默认在 8083 端口提供 REST API,执行:
curl http://localhost:8083/connectors如果各组件均启动成功,将得到如下输出:
[]空数组表示当前没有已注册的 connector,Kafka Connect 服务本身已就绪。
使用 TDengine Sink Connector 同步 Kafka 数据到 TDengine
TDengine Sink Connector 的作用是将指定 topic 的数据同步到 TDengine。用户无需提前创建数据库和超级表:可以手动指定目标数据库名(配置参数connection.database),也可以按一定规则自动生成(配置参数connection.database.prefix)。
从实现原理上看,TDengine Sink Connector 内部使用 TDengine 的无模式(Schemaless)写入接口写入数据,目前支持三种数据格式:InfluxDB Line 协议格式、OpenTSDB Telnet 协议格式和OpenTSDB JSON 协议格式。无模式写入会自动根据实际数据创建超级表、子表并动态扩展列,这正是 Sink Connector 无需预建表结构的底层支撑。
下面的示例将 topicmeters的数据同步到目标数据库power,数据格式为 InfluxDB Line 协议。
添加 Sink Connector 配置文件
mkdir ~/test cd ~/test vi sink-demo.jsonsink-demo.json内容如下:
{ "name": "TDengineSinkConnector", "config": { "connector.class":"com.taosdata.kafka.connect.sink.TDengineSinkConnector", "tasks.max": "1", "topics": "meters", "connection.url": "jdbc:TAOS://127.0.0.1:6030", "connection.user": "root", "connection.password": "taosdata", "connection.database": "power", "db.schemaless": "line", "data.precision": "ns", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.storage.StringConverter", "errors.tolerance": "all", "errors.deadletterqueue.topic.name": "dead_letter_topic", "errors.deadletterqueue.topic.replication.factor": 1 } }关键配置说明:
"topics": "meters"与"connection.database": "power":表示订阅 topicmeters的数据并写入数据库power;"db.schemaless": "line":表示数据使用 InfluxDB Line 协议格式;"data.precision": "ns":声明写入数据的时间戳精度为纳秒,与自动建库的纳秒精度保持一致;errors.tolerance: all与 deadletterqueue 相关配置:允许将解析失败的记录投递到死信 topic,避免单条坏数据阻塞整个消费链路。
创建 Sink Connector 实例
通过 Kafka Connect REST API 提交配置:
curl -X POST -d @sink-demo.json http://localhost:8083/connectors -H "Content-Type: application/json"若命令执行成功,将返回如下内容:
{ "name": "TDengineSinkConnector", "config": { "connection.database": "power", "connection.password": "taosdata", "connection.url": "jdbc:TAOS://127.0.0.1:6030", "connection.user": "root", "connector.class": "com.taosdata.kafka.connect.sink.TDengineSinkConnector", "data.precision": "ns", "db.schemaless": "line", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "tasks.max": "1", "topics": "meters", "value.converter": "org.apache.kafka.connect.storage.StringConverter", "name": "TDengineSinkConnector", "errors.tolerance": "all", "errors.deadletterqueue.topic.name": "dead_letter_topic", "errors.deadletterqueue.topic.replication.factor": "1", }, "tasks": [], "type": "sink" }返回体中"type": "sink"表示该 connector 以 Sink 类型注册成功,随后 Kafka Connect 会自动调度 task 开始消费 topic 并写入 TDengine。
写入测试数据
准备测试数据文本文件test-data.txt,内容如下(每行均为一条 InfluxDB Line 协议记录,最后一段为纳秒时间戳):
meters,location=California.LosAngeles,groupid=2 current=11.8,voltage=221,phase=0.28 1648432611249000000 meters,location=California.LosAngeles,groupid=2 current=13.4,voltage=223,phase=0.29 1648432611250000000 meters,location=California.LosAngeles,groupid=3 current=10.8,voltage=223,phase=0.29 1648432611249000000 meters,location=California.LosAngeles,groupid=3 current=11.3,voltage=221,phase=0.35 1648432611250000000使用kafka-console-producer向 topicmeters灌入测试数据:
cat test-data.txt | kafka-console-producer.sh --broker-list localhost:9092 --topic meters:::note 如果目标数据库power不存在,TDengine Sink Connector 会自动创建数据库。自动建库使用的时间精度为纳秒,这就要求写入数据的时间戳精度也必须是纳秒;如果写入数据的时间戳精度不是纳秒,将抛出异常。 :::
验证同步是否成功
使用taos命令行工具验证:
taos> use power; Database changed. taos> select * from meters; _ts | current | voltage | phase | groupid | location | =============================================================================================================================================================== 2022-03-28 09:56:51.249000000 | 11.800000000 | 221.000000000 | 0.280000000 | 2 | California.LosAngeles | 2022-03-28 09:56:51.250000000 | 13.400000000 | 223.000000000 | 0.290000000 | 2 | California.LosAngeles | 2022-03-28 09:56:51.249000000 | 10.800000000 | 223.000000000 | 0.290000000 | 3 | California.LosAngeles | 2022-03-28 09:56:51.250000000 | 11.300000000 | 221.000000000 | 0.350000000 | 3 | California.LosAngeles | Query OK, 4 row(s) in set (0.004208s)若查询到上述数据,说明同步成功。若未查询到数据,请检查 Kafka Connect 的日志,并结合下文配置参考核对参数。可以观察到:Line 协议中的measurement段(meters)被映射为超级表名,tag_set(location、groupid)被映射为标签列,field_set(current、voltage、phase)被映射为普通数据列——这正是无模式写入的"measurement 即超级表、tag 即标签、field 即列"映射规则在 Sink 场景中的直接体现,详见无模式写入文档。
使用 TDengine Source Connector 同步 TDengine 数据到 Kafka
TDengine Source Connector 的作用是将 TDengine 某个数据库在某一时刻之后的数据推送到 Kafka。其实现原理是:先分批拉取历史数据,再用定时查询的策略同步增量数据;同时会监控表的变化,可以自动同步新增的表。如果 Kafka Connect 重启,一般会从上次中断的位置继续同步,保证断点续传。
从数据格式上看,TDengine Source Connector 会将 TDengine 数据表中的数据转换为InfluxDB Line 协议格式或OpenTSDB JSON 协议格式,再写入 Kafka。
下面示例将数据库test中的数据同步到 topictdengine-test-meters。
添加 Source Connector 配置文件
vi source-demo.json输入以下内容:
{ "name":"TDengineSourceConnector", "config":{ "connector.class": "com.taosdata.kafka.connect.source.TDengineSourceConnector", "tasks.max": 1, "subscription.group.id": "source-demo", "connection.url": "jdbc:TAOS://127.0.0.1:6030", "connection.user": "root", "connection.password": "taosdata", "connection.database": "test", "connection.attempts": 3, "connection.backoff.ms": 5000, "topic.prefix": "tdengine", "topic.delimiter": "-", "poll.interval.ms": 1000, "fetch.max.rows": 100, "topic.per.stable": true, "topic.ignore.db": false, "out.format": "line", "data.precision": "ms", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.storage.StringConverter" } }关键配置说明:
"topic.per.stable": true且"topic.ignore.db": false:表示一个超级表对应一个 Kafka topic,topic 命名规则为<topic.prefix><topic.delimiter><connection.database><topic.delimiter><stable.name>。结合本示例topic.prefix=tdengine、topic.delimiter=-、connection.database=test、超级表名为meters,生成的 topic 即为tdengine-test-meters;"subscription.group.id": "source-demo":指定 TDengine 数据订阅的消费组 ID,Source Connector 默认采用订阅方式(read.method默认为subscription)读取增量数据;"out.format": "line":输出格式为 InfluxDB Line 协议;"data.precision": "ms":时间戳精度为毫秒。
准备测试数据
准备生成测试数据的 SQL 文件prepare-source-data.sql:
DROP DATABASE IF EXISTS test; CREATE DATABASE test; USE test; CREATE STABLE meters (ts TIMESTAMP, current FLOAT, voltage INT, phase FLOAT) TAGS (location BINARY(64), groupId INT); INSERT INTO d1001 USING meters TAGS('California.SanFrancisco', 2) VALUES('2018-10-03 14:38:05.000',10.30000,219,0.31000) \ d1001 USING meters TAGS('California.SanFrancisco', 2) VALUES('2018-10-03 14:38:15.000',12.60000,218,0.33000) \ d1001 USING meters TAGS('California.SanFrancisco', 2) VALUES('2018-10-03 14:38:16.800',12.30000,221,0.31000) \ d1002 USING meters TAGS('California.SanFrancisco', 3) VALUES('2018-10-03 14:38:16.650',10.30000,218,0.25000) \ d1003 USING meters TAGS('California.LosAngeles', 2) VALUES('2018-10-03 14:38:05.500',11.80000,221,0.28000) \ d1003 USING meters TAGS('California.LosAngeles', 2) VALUES('2018-10-03 14:38:16.600',13.40000,223,0.29000) \ d1004 USING meters TAGS('California.LosAngeles', 3) VALUES('2018-10-03 14:38:05.000',10.80000,223,0.29000) \ d1004 USING meters TAGS('California.LosAngeles', 3) VALUES('2018-10-03 14:38:06.500',11.50000,221,0.35000);使用taos命令行执行 SQL 文件:
taos -f prepare-source-data.sql创建 Source Connector 实例
curl -X POST -d @source-demo.json http://localhost:8083/connectors -H "Content-Type: application/json"查看 topic 数据
使用kafka-console-consumer监控 topictdengine-test-meters中的数据:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --from-beginning --topic tdengine-test-meters启动后,一开始会输出所有历史数据,输出格式为 InfluxDB Line 协议:
...... meters,location="California.SanFrancisco",groupid=2i32 current=10.3f32,voltage=219i32,phase=0.31f32 1538548685000000000 meters,location="California.SanFrancisco",groupid=2i32 current=12.6f32,voltage=218i32,phase=0.33f32 1538548695000000000 ......此时会显示全部历史数据。切换到taosshell,插入两条新数据:
USE test; INSERT INTO d1001 VALUES (now, 13.3, 229, 0.38); INSERT INTO d1002 VALUES (now, 16.3, 233, 0.22);切回kafka-console-consumer窗口,可以看到刚插入的 2 条数据已被即时打印,证明增量同步生效。可以观察到输出中的字段均带类型后缀(如10.3f32、219i32、0.31f32),这与无模式写入协议中"数值类型用后缀区分"的规则一一对应(f32 为 float、i32 为 int、无后缀为 double),使 Kafka 侧的数据可以直接被 Sink Connector 或其他无模式写入方消费。
卸载插件
测试完毕后,使用 DELETE 请求停止已加载的 connector。先查看当前活跃的 connector:
curl http://localhost:8083/connectors如果按照前述操作,此时应有两个活跃的 connector。使用以下命令卸载:
curl -X DELETE http://localhost:8083/connectors/TDengineSinkConnector curl -X DELETE http://localhost:8083/connectors/TDengineSourceConnector性能调优
如果在从 TDengine 同步数据到 Kafka 的过程中发现性能不达预期,可以打开$KAFKA_HOME/config/producer.properties配置文件,按以下参数调整 Kafka 生产者的写入吞吐量:
| 参数 | 说明 | 设置建议 |
|---|---|---|
| producer.type | 设置消息发送方式,默认值为sync(同步发送),async表示异步发送。采用异步发送能够提升消息发送的吞吐量。 | async |
| request.required.acks | 配置生产者发送消息后需要等待的确认数量。设置为 1 时,只要领导者副本成功写入消息即向生产者发送确认,无需等待集群中其他副本写入成功。该设置能在一定程度上保证消息可靠性,同时保证吞吐量,因为不需要等待所有副本写入成功,可以减少生产者等待时间,提高发送效率。 | 1 |
| max.request.size | 决定生产者单次请求中可以发送的最大数据量,默认值为 1048576(1M)。设置过小会导致频繁的网络请求、降低吞吐量;设置过大会导致内存占用过高,或在网络状况不佳时增加请求失败概率。建议设置为 100M。 | 104857600 |
| batch.size | 设定 batch 的大小,默认值为 16384(16KB)。消息发送过程中,发送到 Kafka 缓冲区中的消息会被划分成一个个 batch。减小 batch 有助于降低消息延迟,增大 batch 有利于提升吞吐量,可根据实际数据量合理配置,建议设置为 512K。 | 524288 |
| buffer.memory | 设置生产者缓冲待发送消息的内存总量。较大的缓冲区允许生产者积累更多消息后批量发送,提高吞吐量,但也会增加延迟和内存占用。可根据机器资源配置,建议设置为 1G。 | 1073741824 |
这些参数主要作用于 Kafka 生产者端:异步发送(async)避免逐条同步等待,适度放宽batch.size与buffer.memory让更多消息在内存中聚合成批后再发送,配合acks=1在可靠性与吞吐之间取得平衡,适用于 TDengine 大批量历史数据持续入湖的同步场景。
配置参考
通用配置
以下配置项对 TDengine Sink Connector 和 TDengine Source Connector 均适用:
name:connector 名称。connector.class:connector 的完整类名,例如com.taosdata.kafka.connect.sink.TDengineSinkConnector。tasks.max:最大任务数,默认 1。topics:需要同步的 topic 列表,多个用逗号分隔,如topic1,topic2。connection.url:TDengine JDBC 连接字符串,如jdbc:TAOS://127.0.0.1:6030。connection.user:TDengine 用户名,默认root。connection.password:TDengine 用户密码,默认taosdata。connection.attempts:最大尝试连接次数,默认 3。connection.backoff.ms:创建连接失败后的重试间隔时间,单位为毫秒,默认 5000。data.precision:使用 InfluxDB 行协议格式时时间戳的精度,可选值:ms:毫秒;us:微秒;ns:纳秒。
TDengine Sink Connector 特有的配置
connection.database:目标数据库名。如果指定的数据库不存在则自动创建,自动建库使用的时间精度为纳秒。默认值为null;为null时目标数据库命名规则参考connection.database.prefix参数。connection.database.prefix:当connection.database为null时目标数据库的前缀,可以包含占位符${topic}。例如kafka_${topic},对于 topicorders将写入数据库kafka_orders。默认null;为null时目标数据库名与 topic 名一致。batch.size:分批写入时每批的记录数。当 Sink Connector 一次接收到的数据大于该值时将分批写入。max.retries:发生错误时的最大重试次数,默认 1。retry.backoff.ms:发送错误时重试的时间间隔,单位毫秒,默认 3000。db.schemaless:数据格式,可选值:line:InfluxDB 行协议格式;json:OpenTSDB JSON 格式;telnet:OpenTSDB Telnet 行协议格式。
TDengine Source Connector 特有的配置
connection.database:源数据库名称,无缺省值,必须显式指定。topic.prefix:数据导入 Kafka 时使用的 topic 名称前缀,默认为空字符串""。timestamp.initial:数据同步起始时间,格式为yyyy-MM-dd HH:mm:ss;若未指定则从指定 DB 中最早的一条记录开始。poll.interval.ms:检查是否有新建或删除表的时间间隔,单位毫秒,默认 1000。fetch.max.rows:检索数据库时单次最大检索条数,默认 100。query.interval.ms:从 TDengine 一次读取数据的时间跨度,需要根据表中的数据特征合理配置,避免单次查询数据量过大或过小;建议在具体环境中通过测试设置一个较优值,默认值为 0,即获取到当前最新时间的所有数据。out.format:结果集输出格式。line表示输出 InfluxDB Line 协议格式,json表示输出 JSON 格式,默认为line。topic.per.stable:若为true,表示一个超级表对应一个 Kafka topic,topic 命名规则为<topic.prefix><topic.delimiter><connection.database><topic.delimiter><stable.name>;若为false,则指定 DB 中的所有数据进入一个 Kafka topic,topic 命名规则为<topic.prefix><topic.delimiter><connection.database>。topic.ignore.db:topic 命名规则是否包含 database 名称。true表示规则为<topic.prefix><topic.delimiter><stable.name>;false表示规则为<topic.prefix><topic.delimiter><connection.database><topic.delimiter><stable.name>,默认false。此配置项在topic.per.stable设置为false时不生效。topic.delimiter:topic 名称分割符,默认为-。read.method:从 TDengine 读取数据的方式,query或subscription,默认为subscription。subscription.group.id:指定 TDengine 数据订阅的组 ID,当read.method为subscription时此项为必填项。subscription.from:指定 TDengine 数据订阅起始位置,latest或earliest,默认为latest。
原理纵深:Source Connector 背后的数据订阅机制
TDengine Source Connector 默认以subscription方式读取数据,其底层正是 TDengine 内置的数据订阅能力。TDengine 的数据订阅提供与消息队列产品类似的接口:用户在 TDengine 中定义 topic(topic 可以是一个数据库、一张超级表,或对现有表的查询语句),消费者订阅后即可实时收到最新写入的数据。TDengine 会自动索引预写日志(WAL)文件以实现快速随机访问,并提供文件轮转与保留策略,将 WAL 变成持久化、保序的存储引擎,从而支撑"订阅 + ACK + 断点续传"的消费语义。
Source Connector 的"先批量拉历史、再订阅增量、自动感知新表、重启续传"行为,正是对这套订阅机制的封装:多个消费者还可以组成消费组共享消费进度,实现多线程/分布式消费。理解了这一层,就能明白为什么subscription.group.id在订阅模式下是必填项,以及为什么 Connector 重启后能从上一次中断的位置继续同步——这些语义均由 TDengine 侧的消费组与 ACK 机制提供保证。
补充说明
- 除本教程的 Kafka Connect 插件方案外,TDengine 企业版还可在 taosExplorer 中通过可视化界面配置零代码 Kafka 数据写入与数据发布,相关操作可参考零代码数据写入 · Kafka与数据发布 · Kafka;基于主题的消费语义也可对照数据订阅文档进一步理解。
- 本文所有示例均在 Kafka 分布式(connect-distributed)模式下运行;关于如何在独立(standalone)Kafka 环境中使用 Kafka Connect 插件,可参考 Apache Kafka 官方 Connect 文档。
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考