
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/TDengineTDengine Kafka Connector 是面向 Kafka Connect 生态的一对插件TDengine Source Connector 与 TDengine Sink Connector只需一份简单的 JSON 配置即可实现 Kafka 指定 topic 与 TDengine 指定数据库之间的批量或实时双向同步。本文将以官方文档为主线结合仓库中的无模式Schemaless写入文档与数据订阅文档完整演示从环境搭建、插件编译安装、Sink/Source 两个方向的端到端实战并给出全部配置参数的参考说明帮助你在真实 IIoT 场景中快速搭建 Kafka 与 TDengine 之间的数据管道。什么是 Kafka ConnectKafka Connect 是 Apache Kafka 的一个组件用于让其他系统数据库、云服务、文件系统等方便地接入 Kafka。数据既可以通过 Kafka Connect 从外部系统流入 Kafka也可以从 Kafka 流向外部系统。其中Source Connector从其他系统读取数据交给 Kafka Connect 写入 KafkaSink Connector从 Kafka Connect 接收 Kafka 中的数据写入其他系统。需要注意的是Source Connector 和 Sink Connector 都不会直接连接 Kafka Broker——Source Connector 把数据转交给 Kafka ConnectSink Connector 从 Kafka Connect 接收数据。TDengine Kafka Connector 基于这一模型实现了两个插件TDengine Source Connector实时从 TDengine 读取数据并发送给 Kafka ConnectTDengine 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_PKGkafka_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_profileexport 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.skiptrue 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.pathplugin.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如果各组件均启动成功将得到如下输出[]空数组表示当前没有已注册的 connectorKafka Connect 服务本身已就绪。使用 TDengine Sink Connector 同步 Kafka 数据到 TDengineTDengine 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的数据并写入数据库powerdb.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,locationCalifornia.LosAngeles,groupid2 current11.8,voltage221,phase0.28 1648432611249000000 meters,locationCalifornia.LosAngeles,groupid2 current13.4,voltage223,phase0.29 1648432611250000000 meters,locationCalifornia.LosAngeles,groupid3 current10.8,voltage223,phase0.29 1648432611249000000 meters,locationCalifornia.LosAngeles,groupid3 current11.3,voltage221,phase0.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_setlocation、groupid被映射为标签列field_setcurrent、voltage、phase被映射为普通数据列——这正是无模式写入的measurement 即超级表、tag 即标签、field 即列映射规则在 Sink 场景中的直接体现详见无模式写入文档。使用 TDengine Source Connector 同步 TDengine 数据到 KafkaTDengine 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 topictopic 命名规则为topic.prefixtopic.delimiterconnection.databasetopic.delimiterstable.name。结合本示例topic.prefixtdengine、topic.delimiter-、connection.databasetest、超级表名为meters生成的 topic 即为tdengine-test-meterssubscription.group.id: source-demo指定 TDengine 数据订阅的消费组 IDSource Connector 默认采用订阅方式read.method默认为subscription读取增量数据out.format: line输出格式为 InfluxDB Line 协议data.precision: ms时间戳精度为毫秒。准备测试数据准备生成测试数据的 SQL 文件prepare-source-data.sqlDROP 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,locationCalifornia.SanFrancisco,groupid2i32 current10.3f32,voltage219i32,phase0.31f32 1538548685000000000 meters,locationCalifornia.SanFrancisco,groupid2i32 current12.6f32,voltage218i32,phase0.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。先查看当前活跃的 connectorcurl 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表示异步发送。采用异步发送能够提升消息发送的吞吐量。asyncrequest.required.acks配置生产者发送消息后需要等待的确认数量。设置为 1 时只要领导者副本成功写入消息即向生产者发送确认无需等待集群中其他副本写入成功。该设置能在一定程度上保证消息可靠性同时保证吞吐量因为不需要等待所有副本写入成功可以减少生产者等待时间提高发送效率。1max.request.size决定生产者单次请求中可以发送的最大数据量默认值为 10485761M。设置过小会导致频繁的网络请求、降低吞吐量设置过大会导致内存占用过高或在网络状况不佳时增加请求失败概率。建议设置为 100M。104857600batch.size设定 batch 的大小默认值为 1638416KB。消息发送过程中发送到 Kafka 缓冲区中的消息会被划分成一个个 batch。减小 batch 有助于降低消息延迟增大 batch 有利于提升吞吐量可根据实际数据量合理配置建议设置为 512K。524288buffer.memory设置生产者缓冲待发送消息的内存总量。较大的缓冲区允许生产者积累更多消息后批量发送提高吞吐量但也会增加延迟和内存占用。可根据机器资源配置建议设置为 1G。1073741824这些参数主要作用于 Kafka 生产者端异步发送async避免逐条同步等待适度放宽batch.size与buffer.memory让更多消息在内存中聚合成批后再发送配合acks1在可靠性与吞吐之间取得平衡适用于 TDengine 大批量历史数据持续入湖的同步场景。配置参考通用配置以下配置项对 TDengine Sink Connector 和 TDengine Source Connector 均适用nameconnector 名称。connector.classconnector 的完整类名例如com.taosdata.kafka.connect.sink.TDengineSinkConnector。tasks.max最大任务数默认 1。topics需要同步的 topic 列表多个用逗号分隔如topic1,topic2。connection.urlTDengine JDBC 连接字符串如jdbc:TAOS://127.0.0.1:6030。connection.userTDengine 用户名默认root。connection.passwordTDengine 用户密码默认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数据格式可选值lineInfluxDB 行协议格式jsonOpenTSDB JSON 格式telnetOpenTSDB 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 topictopic 命名规则为topic.prefixtopic.delimiterconnection.databasetopic.delimiterstable.name若为false则指定 DB 中的所有数据进入一个 Kafka topictopic 命名规则为topic.prefixtopic.delimiterconnection.database。topic.ignore.dbtopic 命名规则是否包含 database 名称。true表示规则为topic.prefixtopic.delimiterstable.namefalse表示规则为topic.prefixtopic.delimiterconnection.databasetopic.delimiterstable.name默认false。此配置项在topic.per.stable设置为false时不生效。topic.delimitertopic 名称分割符默认为-。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 中定义 topictopic 可以是一个数据库、一张超级表或对现有表的查询语句消费者订阅后即可实时收到最新写入的数据。TDengine 会自动索引预写日志WAL文件以实现快速随机访问并提供文件轮转与保留策略将 WAL 变成持久化、保序的存储引擎从而支撑订阅 ACK 断点续传的消费语义。Source Connector 的先批量拉历史、再订阅增量、自动感知新表、重启续传行为正是对这套订阅机制的封装多个消费者还可以组成消费组共享消费进度实现多线程/分布式消费。理解了这一层就能明白为什么subscription.group.id在订阅模式下是必填项以及为什么 Connector 重启后能从上一次中断的位置继续同步——这些语义均由 TDengine 侧的消费组与 ACK 机制提供保证。补充说明除本教程的 Kafka Connect 插件方案外TDengine 企业版还可在 taosExplorer 中通过可视化界面配置零代码 Kafka 数据写入与数据发布相关操作可参考零代码数据写入 · Kafka与数据发布 · Kafka基于主题的消费语义也可对照数据订阅文档进一步理解。本文所有示例均在 Kafka 分布式connect-distributed模式下运行关于如何在独立standaloneKafka 环境中使用 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),仅供参考