SeaTunnel TDengine 连接器演进全解析:从 Source/Sink 能力到 2.3.12 版本变更实战指南
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
SeaTunnel 的 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 source | 2.3.12 |
| [improve] tdengine options | 2.3.12 |
| [Feature][Connector-V2] Support multi-table sink feature for TDengine | 2.3.11 |
| [Feature][Checkpoint] Add check script for source/sink state class serialVersionUID missing | 2.3.11 |
| [Fix][Connector-V2] Fix NullPointerException when column or tag contains null value in TDengine sink | 2.3.11 |
| [Fix][Connector][TDEngine] TDEngine support NCHAR type | 2.3.9 |
| [Improve][dist] add shade check rule | 2.3.9 |
| [Feature][Restapi] Allow metrics information to be associated to logical plan nodes | 2.3.9 |
| [Improve][Connector-V2] Close all ResultSet after used | 2.3.8 |
| [Fix][Connector-tdengine] Fix sql exception and concurrentmodifyexception when connect to taos and read data | 2.3.7 |
| [Bugfix][TDengine] Fix the issue of losing the driver due to multiple calls to the submit job REST API | 2.3.5 |
| [improve][connector-tdengine] support read bool column from tdengine | 2.3.4 |
| [Bugfix][TDengine] Fix the degree of multiple parallelism affects driver loading | 2.3.4 |
| [Improve][Common] Introduce new error define rule | 2.3.4 |
[Improve] Remove useSeaTunnelSink::getConsumedTypemethod and mark it as deprecated | 2.3.4 |
| [Improve][CheckStyle] Remove useless 'SuppressWarnings' annotation of checkstyle | 2.3.4 |
| [Hotfix][Connector] Fixed TDengine connector using jdbc driver to cause loading error | 2.3.2 |
| [Improve][build] Give the maven module a human readable name | 2.3.1 |
| [Improve][Project] Code format with spotless plugin | 2.3.1 |
| [Feature][Connector-V2] add tdengine source | 2.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)、空值 NPE(2.3.11)。
- 工程规范主线:模块命名、spotless 格式化、shade 检查、错误码规范、ResultSet 资源释放、序列化版本号检查等。
2.3.1:TDengine 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:
| 参数 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| url | String | 是 | - | TDengine REST JDBC URL,格式jdbc:TAOS-RS://host:port,例如jdbc:TAOS-RS://localhost:6041/ |
| username | String | 是 | - | TDengine 认证用户名 |
| password | String | 是 | - | TDengine 认证密码 |
| database | String | 是 | - | TDengine 数据库名,服务端必须已存在 |
| stable | String | 是 | - | 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.4:BOOL 列读取与错误规范引入
支持读取 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.7~2.3.8:并发稳定性与资源释放
并发读取异常修复
2.3.7 修复了连接 TAOS 读取数据时的 SQL 异常与ConcurrentModificationException。这对应源码中的分片分配并发模型:TDengineSourceSplitEnumerator内部使用ConcurrentHashMap维护 pending splits,配合stateLock保证run()、addSplitsBack()、registerReader()等并发入口对分片状态的修改是安全的,详见 TDengineSourceSplitEnumerator.java。
统一释放 ResultSet
2.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.9:NCHAR 类型与发布工程改进
此版本的两个关键点:
- 支持 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_tables | List | 否 | - | 要读取的子表名列表,不配置则读取超级表下全部子表;配置后只读列出的子表,名称必须与服务端完全一致 |
| read_columns | List | 否 | - | 要读取的字段列表,不配置则读取全部列;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_columns,Source 输出的第一列永远是保留字段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 / BIT | BOOLEAN | 2.3.4 起支持 |
| TINYINT / SMALLINT / MEDIUMINT / INT / INTEGER / YEAR | INT | 含对应的 UNSIGNED(INT UNSIGNED 除外) |
| INT UNSIGNED / INTEGER UNSIGNED / BIGINT | LONG | |
| BIGINT UNSIGNED | DECIMAL(20, 0) | |
| DECIMAL | DECIMAL(38, 18) | 源码会打印溢出告警日志 |
| DECIMAL UNSIGNED | DECIMAL(38, 18) | |
| FLOAT | FLOAT | FLOAT UNSIGNED 有溢出告警 |
| DOUBLE | DOUBLE | DOUBLE UNSIGNED 有溢出告警 |
| CHAR / NCHAR / VARCHAR / TEXT 系列 / JSON | STRING | NCHAR 为 2.3.9 起支持 |
| DATE | LOCAL_DATE | |
| TIME | LOCAL_TIME | |
| DATETIME / TIMESTAMP | LOCAL_DATE_TIME | |
| BLOB / BINARY / VARBINARY 系列 | BYTES | |
| GEOMETRY / 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开头。
Sink:exactly-once ✅、CDC 写入 ❌、多表写入 ✅、定时 flush ❌。输入行必须满足"子表名 + 普通列 + TAG 值"的超级表写入形态。
版本升级建议(结合变更日志)
- 从 2.3.1~2.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),仅供参考