SeaTunnel TDengine Sink 连接器完全指南:从超级表写入到多表分发
发布时间:2026/9/20 1:48:49 锦皓数字建站

数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载本篇指南聚焦 SeaTunnel 中的 TDengine 数据接收器Sink讲解如何通过TDengine插件将 SeaTunnel 的批式或流式数据写入 TDengine 数据库的超级表STable。读完本文你将掌握连接器的全部配置参数、输入数据必须满足的子表名 普通列 TAGS写入结构、单表与多表两种写入场景的完整配置样例以及其底层基于 REST JDBC 的 SQL 组装与元数据探测原理可直接复刻到自己的数据集成作业中。支持的引擎TDengine Sink 连接器在以下 SeaTunnel 引擎上均可运行SparkFlinkSeaTunnel Zeta概述与核心机制TDengine 是面向时序数据的高性能数据库其核心数据模型由数据库Database、超级表STable与子表Sub-table构成超级表定义了普通列measurement 列与 TAGS 标签列的公共 Schema子表则通过具体 TAGS 值区分。SeaTunnel 的 TDengine Sink 正是围绕这一模型设计的——它将 SeaTunnel 行数据以一行一条子表记录的方式写入目标超级表。运行 SeaTunnel 任务前需要先创建目标数据库和超级表。该 Sink 支持单表写入也支持在stable中使用${table_name}这类占位符完成多表写入。从源码结构看该连接器的实现位于seatunnel-connectors-v2/connector-tdengine/核心链路为TDengineSinkFactory.java负责插件注册工厂标识为TDengine、参数校验与 Sink 实例创建TDengineSink.java实现SupportMultiTableSink声明连接器具备多表 Sink 能力TDengineSinkWriter.java负责 JDBC 连接、超级表元数据探测与逐行组装INSERT语句。主要特性根据 connector-v2-features.md 中的能力清单TDengine Sink 的支持情况如下精确一次Exactly-once变更数据捕获CDC支持多表写入Multi-table Sink定时刷新Scheduled refresh其中多表写入能力由 TDengineSink.java 中实现的SupportMultiTableSink接口提供。选项详解名称类型是否必传默认值说明urlString是-TDengine REST JDBC 连接地址例如jdbc:TAOS-RS://localhost:6041/。usernameString是-连接 TDengine 使用的用户名。passwordString是-连接 TDengine 使用的密码。databaseString是-TDengine 数据库名称。stableString是-TDengine 超级表名称。多表写入时可以使用占位符例如${table_name}。timezoneString否UTCTDengine 服务端时区用于时间戳转换。write_columnsList否-要写入 TDengine 的普通列名列表。不配置时按目标超级表的列顺序写入不要包含子表名列或 TAGS 字段。common-options否-Sink 插件通用参数请参考 Sink Common Options。上述参数的必传约束在 TDengineSinkFactory.java 的OptionRule中定义url、username、password、database、stable为 requiredtimezone与multi_table_sink_replica为 optional。url [String]TDengine REST JDBC 连接地址对应源码中 TDengineCommonOptions.URL 的url键格式为jdbc:TAOS-RS://host:port。例如jdbc:TAOS-RS://localhost:6041/在 TDengineSinkWriter 中该地址会与database、username、password拼接成完整 JDBC URL 后通过DriverManager.getConnection建立连接随后调用checkDriverExist检查 TAOS 驱动是否存在。username [String]连接 TDengine 使用的用户名例如默认管理员root。password [String]连接 TDengine 使用的密码例如taosdata。database [String]TDengine 数据库名称数据库必须已经在服务端存在。连接器会将其直接拼入 JDBC URL也会在元数据探测 SQLdesc database.stable中使用。stable [String]TDengine 超级表名称。这里有一个需要特别澄清的行为差异该值会被 Sink Writer 原样使用TDengine 连接器本身不会执行占位符替换因此${table_name}等占位符不会在运行时被逐行替换。在多表写入场景下上游 SeaTunnel 框架TablePlaceholderProcessor可能在任务初始化阶段根据上游CatalogTable的标识替换一次stable但这取决于上游框架的接线并不是 TDengine 连接器自身的能力。从 TDengineSinkWriter.java 可以看到Writer 构造时执行的是desc database.stable原样 SQL若stable未被替换为实际超级表名该语句将无法返回有效元数据。timezone [String]TDengine 服务端时区用于时间戳转换默认值为UTC。如果服务端不是 UTC 时区请把该项设置为与服务端一致的时区。对应 TDengineSinkOptions.TIMEZONE 中的默认值定义。在 TDengineSinkWriter.convertDataType 中LocalDateTime类型字段会先按系统默认时区定位再通过withZoneSameInstant(ZoneId.of(config.getTimezone()))转换到目标时区最后格式化为yyyy-MM-dd HH:mm:ss.SSS并加上单引号。write_columns [List]要写入 TDengine 的普通列名列表。不配置时TDengine 会按目标超级表的列顺序写入。这里不要包含第一列子表名也不要包含 TAGS 字段连接器会自动从输入数据末尾取出 TAGS 值。从实现细节看TDengineSinkConfig.of 会将配置中的 List 通过String.join(,, ...)拼成逗号分隔字符串再由 TDengineSinkWriter 拼入 SQL 的列名括号中( col1, col2, ... )若未配置则生成空的列名片段。通用选项Sink 插件通用参数如plugin_input、parallelism等请参考 Sink Common Options。多表写入时可以配合通用参数中的multi_table_sink_replica使用——该参数在 TDengineSinkFactory.java 中已被显式声明为可选项用于为每个下游表配置写入副本数从而放大多表写入的并行度。输入数据格式连接器要求每行输入数据符合超级表写入结构第一列为目标子表名字符串。若该子表不存在Sink 会按目标超级表的 schema 自动创建。接下来的列为write_columns中声明的普通列未配置时使用目标超级表的列顺序。末尾几列为 TAGS 值。TAGS 字段的数量会从目标超级表的元数据中读取。例如目标超级表有 2 个 TAGS 字段时输入行最后 2 列会作为 TAGS 值第一列会作为子表名。这一结构在 TDengineSinkWriter 中体现得非常直接构造 Writer 时执行desc database.stable元数据查询统计note列为TAG的行数得到tagsNum每次write时从行末尾截取tagsNum个字段作为 TAGS 值tags从下标 1 开始截取到arity - tagsNum之间的字段作为普通列metrics下标 0 即第一列子表名被跳过最终组装 SQLINSERT INTO 子表名 using 超级表名 tags ( tag值 ) ( 普通列名 ) VALUES ( 普通列值 );值得注意的是Writer 对返回值有校验若executeUpdate返回 0 行会抛出SQL_OPERATION_FAILED类型的 TDengineConnectorException 异常提示insert error便于快速定位写入失败的行。示例写入单个超级表以下配置演示将数据写入单个超级表meters2数据列依次为ts、voltage、current、powerenv { parallelism 2 job.mode BATCH } sink { TDengine { url jdbc:TAOS-RS://localhost:6041/ username root password taosdata database power2 stable meters2 timezone UTC write_columns [ts, voltage, current, power] } }多表写入匹配的超级表以下配置通过 FakeSource 构造两张表meters3、meters4的数据并在 Sink 侧使用stable ${table_name}实现多表分发source { FakeSource { plugin_output fake tables_configs [ { schema { table meters3 fields { device_id string event_time timestamp metric1 float metric2 int metric3 float status_flag boolean notes string location_tag string group_tag int } } rows [ { kind INSERT fields [d2001, 2023-04-22T14:38:05, 10.3, 219, 0.31, true, nc, California.SanFrancisco, 2] } ] }, { schema { table meters4 fields { device_id string event_time timestamp metric1 float metric2 int metric3 float status_flag boolean notes string location_tag string group_tag int } } rows [ { kind INSERT fields [d1005, 2023-04-22T14:38:05, 110.3, 219, 0.31, true, nc, California.SanFrancisco, 2] } ] } ] } } sink { TDengine { url jdbc:TAOS-RS://localhost:6041/ username root password taosdata database power2 stable ${table_name} timezone UTC } }这里的${table_name}会被 TDengine Sink Writer 当作普通字符串原样使用连接器不会按行替换它因此本示例只有在任务初始化阶段由上游框架依据CatalogTable的标识替换stable时才能生效。目标超级表必须已经存在并具有匹配的 TAGS 列。从 TDengineSinkFactory.java 可以看到Sink 创建时接收的是CatalogTable而 TDengineSink.java 通过getWriteCatalogTable()将CatalogTable暴露给框架这正是框架侧在初始化阶段执行占位符替换与多表路由所依赖的接线点。数据类型与时区转换的源码级细节TDengineSinkWriterTest.java 中的testConvertDataTypeWithNull用例精确验证了字段级别的类型转换规则null 值原样保留不套引号、不转换LocalDateTime按timezone配置完成时区转换并格式化为yyyy-MM-dd HH:mm:ss.SSS字符串如2023-04-14 15:30:45.000String自动包裹单引号避免拼入 SQL 时出现语法问题Integer / Double 等数值类型保持原值不处理。该测试同时用 Mock 的 JDBC 结果集模拟了desc test_db.test_stable返回 TAG 行的场景验证了tagsNum的统计逻辑。这意味着实际生产配置时务必保证输入行中 TAGS 字段的类型与目标超级表定义一致字符串型 TAGS 会被自动加引号而数值型 TAGS 保持裸值。常见注意事项数据库与超级表需提前建好连接器不做建库建表操作database与stable必须已存在于服务端否则构造 Writer 时的desc元数据查询会失败。输入行结构必须严格对齐第一列子表名 中间普通列 末尾 TAGS 值TAGS 列数由超级表元数据决定行字段顺序错乱会导致写入错误的数据。${table_name}不是连接器能力占位符替换依赖上游框架在任务初始化阶段按CatalogTable标识进行连接器自身始终原样使用stable字符串。时区一致性timezone默认UTC若 TDengine 服务端使用其他时区需显式配置否则时间戳会被错误转换。write_columns的边界只声明普通列不要包含子表名列与 TAGS 字段TAGS 由连接器自动从行尾截取。变更日志连接器的历次变更记录维护在 connector-tdengine.md对应英文版本位于 docs/en/connectors/changelog/connector-tdengine.md发布新版本后建议同步核对变更内容。赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐SeaTunnel HBase Sink 连接器完全指南从单表写入到多表分发的配置与原理SeaTunnel HBase Sink 连接器完全指南从单表写入到多表分发的配置与原理 本文以 docs/zh/connectors/sink/Hbase.数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Hudi Sink 连接器完全指南多表写入、UPSERT 与 CDC 实践SeaTunnel Hudi Sink 连接器完全指南多表写入、UPSERT 与 CDC 实践 本文以 Apache SeaTunnel 仓库中 Hudi 接数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Pulsar Sink 连接器实战指南从多表路由到 EXACTLY_ONCE 事务写入SeaTunnel Pulsar Sink 连接器实战指南从多表路由到 EXACTLY_ONCE 事务写入 SeaTunnel 的 Pulsar Sink 连数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。