news 2026/9/16 17:30:22

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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/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 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

梳理后可以得到三条清晰的演进主线:

  1. 能力扩展主线:2.3.1 引入 Source 基础能力 → 2.3.4 支持读取 BOOL 列 → 2.3.9 支持 NCHAR 类型 → 2.3.11 支持多表写入 Sink → 2.3.12 支持子表过滤(subtable)与列投影(fieldNames)。
  2. 稳定性修复主线:驱动加载问题(2.3.2、2.3.5)、并行度影响驱动加载(2.3.4)、并发读异常(2.3.7)、空值 NPE(2.3.11)。
  3. 工程规范主线:模块命名、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:

参数类型必填默认值说明
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.4:BOOL 列读取与错误规范引入

支持读取 BOOL 列

TDengine 使用BOOL类型存储布尔值。此版本通过类型映射逻辑补齐了对该类型的支持(详见后文类型映射章节),使BOOL列能够正确映射为 SeaTunnel 的BOOLEAN类型。

工程规范沉淀

同一版本还引入了三项基础规范:

  • 错误码规则:新增 TDengineConnectorErrorCode.java,读写失败统一抛出 TDengineConnectorException.java 并携带明确错误码(如READER_OPERATION_FAILEDWRITER_OPERATION_FAILEDSQL_OPERATION_FAILEDUNSUPPORTED_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 类型与发布工程改进

此版本的两个关键点:

  1. 支持 NCHAR 类型:TDengine 的NCHAR(定长多字节字符串)在类型映射中被归入字符串类型族,映射为 SeaTunnel 的STRING。这也是 2.3.11 之前 Source 端能读、Sink 端能写多语言文本数据的前提。
  2. 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_columns,Source 输出的第一列永远是保留字段subtable_name(子表名)。实现上由 addHiddenAttribute() 将该字段插入到字段列表最前面。该字段正是 Sink 端回写时使用的目标子表名,因此"TDengine 读 → TDengine 写"的管道天然依赖这一约定。

源码级原理:Source 分片发现、查询构建与类型映射

分片发现与轮询分配

Source 的并行读取机制分三层:

  1. 元数据发现TDengineSource.getStableMetadata()通过desc db.stable获取时间戳字段名与字段类型,通过select table_name from information_schema.ins_tables where db_name='...' and stable_name='...'获取子表列表。
  2. 分片构建:TDengineSourceSplitEnumerator.createSplitBySubTable() 为每个子表生成查询 SQL,时间范围遵循左闭右开语义:timestamp_field >= lower_bound and timestamp_field < upper_bound
  3. 轮询分配getSplitOwner()assignCount % numReaders将分片轮流分配给各并行 reader,并在addPendingSplit前按 splitId 排序,保证分配顺序可预期。故障恢复时通过snapshotState()保存shouldEnumerate、pending splits 与分配计数,由restoreEnumerator()重建状态。

数据读取与类型转换

TDengineSourceReader 的读取流程为:

  1. open()阶段使用TSDBDriver.PROPERTY_KEY_USER/PROPERTY_KEY_PASSWORD组装连接属性并建立 JDBC 连接;
  2. pollNext()从并发队列取分片,执行分片 SQL,将java.sql.Timestamp转为LocalDateTimebyte[]转为String,其余类型原样透传;
  3. 所有分片消费完毕且收到handleNoMoreSplits()后,调用context.signalNoMoreElement()结束有界读取。

完整类型映射表

TDengineTypeMapper.java 定义了 TDengine 类型到 SeaTunnel 类型的完整映射规则(另有 TDengineTypeMapperTest.java 覆盖验证):

TDengine 类型SeaTunnel 类型备注
BOOL / BITBOOLEAN2.3.4 起支持
TINYINT / SMALLINT / MEDIUMINT / INT / INTEGER / YEARINT含对应的 UNSIGNED(INT UNSIGNED 除外)
INT UNSIGNED / INTEGER UNSIGNED / BIGINTLONG
BIGINT UNSIGNEDDECIMAL(20, 0)
DECIMALDECIMAL(38, 18)源码会打印溢出告警日志
DECIMAL UNSIGNEDDECIMAL(38, 18)
FLOATFLOATFLOAT UNSIGNED 有溢出告警
DOUBLEDOUBLEDOUBLE UNSIGNED 有溢出告警
CHAR / NCHAR / VARCHAR / TEXT 系列 / JSONSTRINGNCHAR 为 2.3.9 起支持
DATELOCAL_DATE
TIMELOCAL_TIME
DATETIME / TIMESTAMPLOCAL_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 列(locationgroupid)必须放在末尾,这样下游 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 输入表分别落到meters3meters4)则使用stable = "${table_name}"占位符,由上游框架在作业初始化时替换一次,详见 Sink 文档 的完整示例。

Sink 写入原理与注意事项

TDengineSinkWriter 的写入过程分三步:

  1. 构造 JDBC URL:将 url、database、username、password 拼接为url + database + "?user=...&password=...",并调用checkDriverExist确保驱动可用;
  2. 获取 TAG 元数据:执行desc db.stable,统计note列为TAG的字段个数tagsNum
  3. 组装写入 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 端目前最完整的形态。

常见坑位清单

  1. 目标超级表必须预先创建:Sink 通过元数据读取 TAG 个数,表不存在会直接失败;
  2. read_columns的 TAG 列必须放末尾:否则 Sink 无法正确切分普通列与 TAG 值;
  3. subtable_name是保留字段:Source 自动注入、Sink 自动消费,配置列清单时不要手动包含;
  4. 时间范围是左闭右开lower_bound闭、upper_bound开,避免边界数据重复或遗漏;
  5. 时区一致性:Source 端时间戳按Timestamp → LocalDateTime转换,Sink 端按timezone转换,两端 TDengine 服务时区与配置需对齐;
  6. 驱动加载:并行度高时驱动注册可能成为隐患,务必使用包含 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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/16 17:30:21

书霸AI官网www.shubaai.com,公众号搜一搜书霸AI写作

一份开题报告&#xff0c;像是论文正式写作前的“施工图”。题目是否清楚、研究问题是否具体、方法是否可行&#xff0c;都会影响后续论文能否顺利展开。截图中的书霸AI开题报告功能&#xff0c;重点解决的正是前期构思难、结构乱、资料不会组织等问题。使用时&#xff0c;可以…

作者头像 李华
网站建设 2026/9/16 17:28:32

MATLAB滑动窗口S值计算在声发射信号分析中的应用

1. 项目背景与核心价值声发射信号分析在工业无损检测、结构健康监测等领域有着广泛应用。传统分析方法往往采用固定时间窗口&#xff0c;难以捕捉信号中的瞬态特征。滑动窗口技术通过动态分割信号&#xff0c;能够更精准地定位异常事件并提取特征参数。其中S值&#xff08;Sign…

作者头像 李华
网站建设 2026/9/16 17:28:23

STM32H743固件加密:C#上位机与AES-GCM实现详解

简介&#xff1a;一套基于STM32H743单片机生成AES加密固件的上位机软件源码包&#xff0c;面向嵌入式开发者与安全固件研究人员&#xff0c;解决固件加密传输、密钥管理及在线升级等需求。压缩包共808个文件&#xff0c;约20.24MB&#xff0c;涵盖C/C源码&#xff08;.h/.c&…

作者头像 李华
网站建设 2026/9/16 17:28:21

OpenClaw 跑 Agent 任务:Key 用 TaoToken 压 Token 开销

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华