资讯详情

SeaTunnel TDengine 连接器演进全解析:从 Source/Sink 能力到 2.3.12 版本变更实战指南

📅 2026/9/16 17:30:28 | 华诺云谱 👁 阅读
SeaTunnel TDengine 连接器演进全解析:从 Source/Sink 能力到 2.3.12 版本变更实战指南
SeaTunnel TDengine 连接器演进全解析从 Source/Sink 能力到 2.3.12 版本变更实战指南【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 TDengine 连接器负责在 TDengine 时序数据库与 SeaTunnel 数据管道之间打通批量读写通道Source 端以超级表super table为读取对象将一张超级表下的全部子表按分片并行读取Sink 端则按 TDengine 超级表写入模型把上游行数据回写。本文以官方变更日志 connector-tdengine.md 为骨架完整梳理该连接器自 2.3.1 引入以来的每一次关键演进子表过滤、列投影、多表写入、NCHAR/BOOL 类型支持、驱动加载修复等并结合 Source 文档、Sink 文档 与seatunnel-connectors-v2/connector-tdengine模块源码带你掌握每个版本新增参数的用法、底层实现原理与实战配置模板。连接器全景TDengine 在 SeaTunnel 中的定位TDengine 连接器位于 seatunnel-connectors-v2/connector-tdengine 模块通过jdbc:TAOS-RS://host:port形式的 REST JDBC 协议与 TDengine 服务端交互支持 Spark、Flink 与 SeaTunnel Zeta 三种引擎。核心设计理念是围绕 TDengine 的超级表模型工作Source 端一个 source split 对应一张子表sub table。连接器会先通过元数据查询发现超级表下的全部子表再为每个子表生成一条带时间范围过滤条件的 SQL并行读取。Sink 端要求输入行遵循超级表写入形态——第一个字段是目标子表名中间字段是普通列最后几个字段是 TAG 值TAG 个数通过查询目标超级表元数据获得。这种设计让TDengine → TDengine的库间同步、以及 TDengine 与其他系统的集成都变得非常直接。下面的演进时间线将展示这一能力是如何一步步打磨出来的。演进时间线从 2.3.1 到 2.3.12 的完整变更记录以下为 connector-tdengine.md 中记录的完整变更历史外部提交链接已省略保留变更说明与版本对应关系变更版本[Feature][connector-tdengine] Support subtable and fieldNames in tdengine source2.3.12[improve] tdengine options2.3.12[Feature][Connector-V2] Support multi-table sink feature for TDengine2.3.11[Feature][Checkpoint] Add check script for source/sink state class serialVersionUID missing2.3.11[Fix][Connector-V2] Fix NullPointerException when column or tag contains null value in TDengine sink2.3.11[Fix][Connector][TDEngine] TDEngine support NCHAR type2.3.9[Improve][dist] add shade check rule2.3.9[Feature][Restapi] Allow metrics information to be associated to logical plan nodes2.3.9[Improve][Connector-V2] Close all ResultSet after used2.3.8[Fix][Connector-tdengine] Fix sql exception and concurrentmodifyexception when connect to taos and read data2.3.7[Bugfix][TDengine] Fix the issue of losing the driver due to multiple calls to the submit job REST API2.3.5[improve][connector-tdengine] support read bool column from tdengine2.3.4[Bugfix][TDengine] Fix the degree of multiple parallelism affects driver loading2.3.4[Improve][Common] Introduce new error define rule2.3.4[Improve] Remove useSeaTunnelSink::getConsumedTypemethod and mark it as deprecated2.3.4[Improve][CheckStyle] Remove useless SuppressWarnings annotation of checkstyle2.3.4[Hotfix][Connector] Fixed TDengine connector using jdbc driver to cause loading error2.3.2[Improve][build] Give the maven module a human readable name2.3.1[Improve][Project] Code format with spotless plugin2.3.1[Feature][Connector-V2] add tdengine source2.3.1梳理后可以得到三条清晰的演进主线能力扩展主线2.3.1 引入 Source 基础能力 → 2.3.4 支持读取 BOOL 列 → 2.3.9 支持 NCHAR 类型 → 2.3.11 支持多表写入 Sink → 2.3.12 支持子表过滤subtable与列投影fieldNames。稳定性修复主线驱动加载问题2.3.2、2.3.5、并行度影响驱动加载2.3.4、并发读异常2.3.7、空值 NPE2.3.11。工程规范主线模块命名、spotless 格式化、shade 检查、错误码规范、ResultSet 资源释放、序列化版本号检查等。2.3.1TDengine Source 的诞生与基础配置引入背景2.3.1 版本首次将 TDengine 引入 SeaTunnel Connector-V2 生态对应提交 add tdengine source。早期实现的定位是批量BATCH有界读取这一点在源码中仍有明确痕迹——TDengineSource.java 的getBoundedness()返回Boundedness.BOUNDED类注释也保留了 TODO未来优化方向是batch → batch stream与单条写入 → 批量写入。基础参数Source 与 Sink 共享的公共选项Source 和 Sink 复用同一组连接参数定义在 TDengineCommonOptions.java参数类型必填默认值说明urlString是-TDengine REST JDBC URL格式jdbc:TAOS-RS://host:port例如jdbc:TAOS-RS://localhost:6041/usernameString是-TDengine 认证用户名passwordString是-TDengine 认证密码databaseString是-TDengine 数据库名服务端必须已存在stableString是-TDengine 超级表名从 TDengineSourceConfig.java 的buildSourceConfig与 TDengineSinkConfig.java 的of可以看出所有参数最终都会映射到上述连接字段并被封装为可序列化的配置对象下发到每个 reader/writer 并行实例。驱动加载的早期痛点连接器依赖 TDengine 官方 JDBC 驱动com.taosdata.jdbc.TSDBDriver。2.3.2 的 Hotfix 解决了使用 JDBC 驱动导致加载错误的问题2.3.5 又修复了多次调用提交作业 REST API 导致驱动丢失的问题。这些 bug 的根源都在于驱动类在并发/多实例场景下被重复加载或加载失败。目前源码中统一通过 TDengineUtil.java 的checkDriverExist(jdbcUrl)工具方法处理先检查驱动是否存在不存在则尝试注册Source 的 open() 与 Sink 的构造函数都会在建立连接前调用它。这也解释了为什么 2.3.4 要专门修复并行度影响驱动加载——每个并行 reader 实例都会走同一套驱动检查与注册逻辑。2.3.4BOOL 列读取与错误规范引入支持读取 BOOL 列TDengine 使用BOOL类型存储布尔值。此版本通过类型映射逻辑补齐了对该类型的支持详见后文类型映射章节使BOOL列能够正确映射为 SeaTunnel 的BOOLEAN类型。工程规范沉淀同一版本还引入了三项基础规范错误码规则新增 TDengineConnectorErrorCode.java读写失败统一抛出 TDengineConnectorException.java 并携带明确错误码如READER_OPERATION_FAILED、WRITER_OPERATION_FAILED、SQL_OPERATION_FAILED、UNSUPPORTED_DATA_TYPE移除getConsumedType废弃方法规范 Sink 接口用法CheckStyle 清理移除无用的SuppressWarnings注解。2.3.72.3.8并发稳定性与资源释放并发读取异常修复2.3.7 修复了连接 TAOS 读取数据时的 SQL 异常与ConcurrentModificationException。这对应源码中的分片分配并发模型TDengineSourceSplitEnumerator内部使用ConcurrentHashMap维护 pending splits配合stateLock保证run()、addSplitsBack()、registerReader()等并发入口对分片状态的修改是安全的详见 TDengineSourceSplitEnumerator.java。统一释放 ResultSet2.3.8 要求所有使用过的 ResultSet 必须关闭。当前 Source 的元数据查询与数据读取均采用try-with-resources写法TDengineSource.java 中用单条try (Connection / Statement / ResultSet...)同时执行desc db.stable拿字段类型与information_schema.ins_tables拿子表列表TDengineSourceReader.read() 用try (Statement / ResultSet)执行分片 SQL 并逐行产出SeaTunnelRow。2.3.9NCHAR 类型与发布工程改进此版本的两个关键点支持 NCHAR 类型TDengine 的NCHAR定长多字节字符串在类型映射中被归入字符串类型族映射为 SeaTunnel 的STRING。这也是 2.3.11 之前 Source 端能读、Sink 端能写多语言文本数据的前提。shade 检查规则与 RestAPI 指标关联属于发布工程dist 打包与监控侧改进与连接器读写逻辑无直接关系但对依赖打包的稳定性有间接贡献。2.3.11多表写入 Sink 与空值 NPE 修复多表写入multi-table sink这是 Sink 端最重要的能力升级。Sink 接口实现了SupportMultiTableSinkWriter见 TDengineSinkWriter.java配合 Sink 文档中的说明其工作机制为stable参数中可以出现${table_name}占位符由 SeaTunnel 上游框架的TablePlaceholderProcessor在作业初始化阶段根据上游CatalogTable标识替换一次注意TDengine 连接器本身不做逐行替换${table_name}对 writer 而言是字面量因此多表写入依赖上游框架在作业构建时的替换行为并非 TDengine 专属能力配合通用 Sink 参数multi_table_sink_replica可控制多表写入的副本数。空值 NPE 修复修复了列或 TAG 含 null 值时抛 NullPointerException的问题。对应 TDengineSinkWriter.convertDataType() 中object null直接返回 null 的防御逻辑同一版本还加入了 source/sink 状态类serialVersionUID缺失的检查脚本保证 checkpoint 状态在版本间可兼容反序列化。2.3.12子表过滤与列投影的最终形态这是 Source 端能力的集大成版本。新增的两个参数定义在 TDengineSourceOptions.java参数类型必填默认值说明sub_tablesList否-要读取的子表名列表不配置则读取超级表下全部子表配置后只读列出的子表名称必须与服务端完全一致read_columnsList否-要读取的字段列表不配置则读取全部列TAG 列必须放在列表末尾且不要包含subtable_name这两个参数在 TDengineSourceConfig.buildSourceConfig() 中被转换为Set类型并在 TDengineSource.getStableMetadata() 中发挥过滤作用read_columns决定desc db.stable返回的字段哪些进入输出 Schema不匹配的字段被continue跳过sub_tables决定information_schema.ins_tables返回的子表哪些生成 source split。输出 Schema 的隐藏字段无论是否配置read_columnsSource 输出的第一列永远是保留字段subtable_name子表名。实现上由 addHiddenAttribute() 将该字段插入到字段列表最前面。该字段正是 Sink 端回写时使用的目标子表名因此TDengine 读 → TDengine 写的管道天然依赖这一约定。源码级原理Source 分片发现、查询构建与类型映射分片发现与轮询分配Source 的并行读取机制分三层元数据发现TDengineSource.getStableMetadata()通过desc db.stable获取时间戳字段名与字段类型通过select table_name from information_schema.ins_tables where db_name... and stable_name...获取子表列表。分片构建TDengineSourceSplitEnumerator.createSplitBySubTable() 为每个子表生成查询 SQL时间范围遵循左闭右开语义timestamp_field lower_bound and timestamp_field upper_bound。轮询分配getSplitOwner()用assignCount % numReaders将分片轮流分配给各并行 reader并在addPendingSplit前按 splitId 排序保证分配顺序可预期。故障恢复时通过snapshotState()保存shouldEnumerate、pending splits 与分配计数由restoreEnumerator()重建状态。数据读取与类型转换TDengineSourceReader 的读取流程为open()阶段使用TSDBDriver.PROPERTY_KEY_USER/PROPERTY_KEY_PASSWORD组装连接属性并建立 JDBC 连接pollNext()从并发队列取分片执行分片 SQL将java.sql.Timestamp转为LocalDateTime、byte[]转为String其余类型原样透传所有分片消费完毕且收到handleNoMoreSplits()后调用context.signalNoMoreElement()结束有界读取。完整类型映射表TDengineTypeMapper.java 定义了 TDengine 类型到 SeaTunnel 类型的完整映射规则另有 TDengineTypeMapperTest.java 覆盖验证TDengine 类型SeaTunnel 类型备注BOOL / BITBOOLEAN2.3.4 起支持TINYINT / SMALLINT / MEDIUMINT / INT / INTEGER / YEARINT含对应的 UNSIGNEDINT UNSIGNED 除外INT UNSIGNED / INTEGER UNSIGNED / BIGINTLONGBIGINT UNSIGNEDDECIMAL(20, 0)DECIMALDECIMAL(38, 18)源码会打印溢出告警日志DECIMAL UNSIGNEDDECIMAL(38, 18)FLOATFLOATFLOAT UNSIGNED 有溢出告警DOUBLEDOUBLEDOUBLE UNSIGNED 有溢出告警CHAR / NCHAR / VARCHAR / TEXT 系列 / JSONSTRINGNCHAR 为 2.3.9 起支持DATELOCAL_DATETIMELOCAL_TIMEDATETIME / TIMESTAMPLOCAL_DATE_TIMEBLOB / BINARY / VARBINARY 系列BYTESGEOMETRY / UNKNOWN不支持抛出UNSUPPORTED_DATA_TYPE实战配置从简单读取到完整同步管道场景一读取超级表全量子表时间范围过滤配置来自 Source 文档env { parallelism 2 job.mode BATCH } source { TDengine { url jdbc:TAOS-RS://localhost:6041/ username root password taosdata database power stable meters lower_bound 2018-10-03 14:38:05.000 upper_bound 2018-10-03 14:38:16.801 plugin_output tdengine_result } }lower_bound是闭区间upper_bound是开区间时间戳格式需与 TDengine 兼容如2018-10-03 14:38:05.000。场景二只读指定子表与指定列2.3.12 能力source { TDengine { url jdbc:TAOS-RS://localhost:6041/ username root password taosdata database power stable meters lower_bound 2018-10-03 14:38:05.000 upper_bound 2018-10-03 14:38:16.801 sub_tables [d1001, d1002] read_columns [ts, current, voltage, phase, off, nc, location, groupid] } }注意read_columns的顺序决定输出字段顺序普通列在前TAG 列location、groupid必须放在末尾这样下游 TDengine Sink 才能正确拆分普通列与 TAG 值同时不要包含subtable_name该字段由连接器自动加为第一列。场景三TDengine 到 TDengine 的库间同步2.3.11 多表写入形态env { parallelism 2 job.mode BATCH } source { TDengine { url jdbc:TAOS-RS://tdengine-src:6041/ username root password taosdata database power stable meters lower_bound 2018-10-03 14:38:05.000 upper_bound 2018-10-03 14:38:16.801 plugin_output tdengine_result } } sink { TDengine { url jdbc:TAOS-RS://tdengine-sink:6041/ username root password taosdata database power2 stable meters2 timezone UTC } }Sink 依据 Source 输出的subtable_name确定目标子表名目标超级表必须已存在Sink 写入前会执行desc db.stable读取 TAG 个数。timezone参数默认UTC用于时间戳转换若 TDengine 服务端不是 UTC 时区需显式指定。场景四单超级表写入与多输入表写入单表写入时可配合write_columns显式指定普通列不包括子表名与 TAG 列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则使用stable ${table_name}占位符由上游框架在作业初始化时替换一次详见 Sink 文档 的完整示例。Sink 写入原理与注意事项TDengineSinkWriter 的写入过程分三步构造 JDBC URL将 url、database、username、password 拼接为url database ?user...password...并调用checkDriverExist确保驱动可用获取 TAG 元数据执行desc db.stable统计note列为TAG的字段个数tagsNum组装写入 SQL从行的末尾切出tagsNum个字段作为 TAG 值行的第 0 个字段作为子表名中间字段作为普通列生成形如INSERT INTO subTable using stable tags ( ... ) [ ( cols ) ] VALUES ( ... );的语句执行。写入过程中的时间戳转换convertDataType会把LocalDateTime先按系统默认时区解释再转换到配置的timezone默认 UTC并格式化为yyyy-MM-dd HH:mm:ss.SSS字符串字段自动加单引号包裹。实战提示write_columns未配置时按目标超级表的列顺序写入TAG 值与子表名永远由连接器自动处理不要出现在write_columns中。测试用例 TDengineSinkWriterTest.java 覆盖了写入 SQL 组装与类型转换逻辑。连接器能力矩阵与最佳实践能力速览来自官方文档Source批量读取 ✅、流式读取 ❌、exactly-once ✅、列投影 ✅、并行读取 ✅、用户自定义分片 ❌。每个 source split 读取一张子表输出 Schema 恒以subtable_name开头。Sinkexactly-once ✅、CDC 写入 ❌、多表写入 ✅、定时 flush ❌。输入行必须满足子表名 普通列 TAG 值的超级表写入形态。版本升级建议结合变更日志从 2.3.12.3.3 升级优先解决驱动加载稳定性问题2.3.2/2.3.5建议至少升级到 2.3.5需要 BOOL 列升级到 2.3.4需要 NCHAR 文本列升级到 2.3.9需要 Sink 多表写入升级到 2.3.11并注意${table_name}依赖上游框架替换需要子表过滤与列投影升级到 2.3.12这是 Source 端目前最完整的形态。常见坑位清单目标超级表必须预先创建Sink 通过元数据读取 TAG 个数表不存在会直接失败read_columns的 TAG 列必须放末尾否则 Sink 无法正确切分普通列与 TAG 值subtable_name是保留字段Source 自动注入、Sink 自动消费配置列清单时不要手动包含时间范围是左闭右开lower_bound闭、upper_bound开避免边界数据重复或遗漏时区一致性Source 端时间戳按Timestamp → LocalDateTime转换Sink 端按timezone转换两端 TDengine 服务时区与配置需对齐驱动加载并行度高时驱动注册可能成为隐患务必使用包含 2.3.5 驱动修复的版本。参考资料本文章核心依据变更日志Source 连接器文档Sink 连接器文档模块源码seatunnel-connectors-v2/connector-tdengine测试用例TDengineTest.java、TDengineSourceReaderTest.java、TDengineSourceSplitEnumeratorTest.java【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。