SeaTunnel TDengine Source Connector 实战指南:配置、原理与数据同步 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本指南以 TDengine 官方文档 为主线结合connector-tdengine模块源码与 e2e 测试系统讲解如何在 SeaTunnel 中通过 TDengine Source 连接器批量读取 TDengine 时序数据。读完你将掌握该连接器的全部配置参数、正确的连接串与 SQL 语义、按子表拆分并行读取的工作原理、类型映射规则以及可复制运行的完整配置示例。连接器概述TDengine Source 连接器用于读取外部 TDengine 数据源的数据Read external data source data through TDengine在 SeaTunnel 中注册的插件名为TDengine。它是标准的 v2 连接器通过AutoService(SeaTunnelSource.class)自动注册声明位于 TDengineSource.java并在 plugin-mapping.properties 中映射为seatunnel.source.TDengine connector-tdengine。从源码结构看该连接器基于 SeaTunnel Source API 实现了SeaTunnelSource、SourceReader、SourceSplitEnumerator三件套组件类职责SourceTDengineSource.java校验配置、建立连接、获取 Stable 元数据Split EnumeratorTDengineSourceSplitEnumerator.java按子表拆分数据、分配分片、管理快照ReaderTDengineSourceReader.java执行查询、转换数据类型、输出SeaTunnelRowKey features✅ batch批式只支持批处理不支持流模式❌ stream流式✅ exactly-once精确一次❌ column projection列投影——但支持查询 SQL通过查询语句本身可以达到投影效果✅ parallelism并行度按子表拆分实现并行读取❌ 用户自定义分片对应的Boundedness.BOUNDED实现见 TDengineSource.java#L95-L97源码注释也明确指出该连接器当前批读、单条写出后续有待优化为流式与批量写出。配置项详解连接器完整的配置项如下表与官方文档保持一致默认值取自 TDengineSourceConfig.java 源码名称类型是否必填默认值说明urlstring是-TDengine 的 JDBC 连接地址usernamestring是-登录用户名passwordstring是-登录密码databasestring是-数据库名stablestring是-超级表stable名lower_boundlong是-迁移时间窗口下界upper_boundlong是-迁移时间窗口上界url [string]TDengine 的连接地址支持 RESTful 与原生两种连接方式。示例jdbc:TAOS-RS://localhost:6041/驱动加载逻辑见 TDengineUtil.java若 URL 以jdbc:TAOS-RS://开头则加载com.taosdata.jdbc.rs.RestfulDrivertaosadapter REST 驱动默认端口 6041否则加载原生驱动com.taosdata.jdbc.TSDBDriver默认端口 6030。连接器依赖的驱动版本为taos-jdbcdriver 3.0.3见 connector-tdengine/pom.xml。username [string]连接 TDengine 时使用的用户名如root。password [string]连接 TDengine 时使用的密码如taosdata。database [string]要读取数据的数据库名如power。注意该配置是必填项连接器在prepare阶段会调用CheckConfigUtil.checkAllExists强制校验url/database/stable/username/password五个配置全部非空否则抛出CONFIG_VALIDATION_FAILED异常见 TDengineSource.java#L81-L92。stable [string]要读取的超级表stable名如meters。TDengine 中子表subtable会继承超级表的 schema读取时实际遍历的是该超级表下的所有子表。lower_bound [long]迁移时间窗口的下界起始时间。它是一个时间字符串如2018-10-03 14:38:05.000虽然配置表类型标注为 long但实际以字符串形式参与 SQL 拼接。upper_bound [long]迁移时间窗口的上界结束时间。同样以时间字符串形式配置。关于时间窗口的语义见 TDengineSourceSplitEnumerator.java#L100-L112 源码窗口为左闭右开区间——下界生成ts lower_bound条件上界生成ts upper_bound条件两者用and连接。也就是说读取范围覆盖[lower_bound, upper_bound)。timezone [string]源码扩展项官方文档配置表中未列出但源码 TDengineSourceConfig.java#L46-L49 中已预留timezone参数默认值为UTC。其含义是jdbc:TAOS-RS中 timezone 参数只作用于 taosadapter 服务端因此该参数代表服务端时区设置。当前配置解析已支持但分片 SQL 构建时尚未使用该字段可推断为预留扩展项。工作原理从元数据发现到子表拆分1. 元数据获取TDengineSource.prepare()阶段通过 getStableMetadata 完成两件事构造 JDBC URLurl database ?user username password password并调用checkDriverExist检查/注册驱动随后建立连接执行desc database.stable获取超级表列结构与首列时间戳字段名查询information_schema.ins_tables元数据表列出该数据库下隶属于该超级表的所有子表名将列结构映射为SeaTunnelRowType并将子表名列subtable_name作为隐藏字段插到 schema 第一位addHiddenAttribute方便下游识别每条数据来自哪个子表。2. 分片与并行TDengineSourceSplitEnumerator负责把整个读取任务拆分为每个子表一个 splitgetAllSplits()遍历元数据发现的所有子表为每个子表生成一条带时间窗口过滤的查询 SQLTDengineSourceSplitEnumerator.java#L79-L88分片分配采用哈希取模策略(splitId.hashCode() Integer.MAX_VALUE) % numReaders将子表分片按并行度均匀打散到各 ReaderTDengineSourceSplitEnumerator.java#L63-L65每个 split 携带一条最终查询语句TDengineSourceSplit由 Reader 直接执行因此并行度越高、子表越多读取吞吐越大。由于子表是天然的并行单元无需用户自定义 split这也是特性表中“parallelism 支持、用户自定义 split 不支持”的源码级原因。3. 数据读取与类型转换TDengineSourceReader在每个分片上执行查询将 JDBC 结果逐行转为SeaTunnelRow行首写入split.splitId()作为子表名对应隐藏列subtable_name类型转换规则见 TDengineSourceReader.java#L153-L160Timestamp转为LocalDateTimebyte[]转为String读取完成后调用signalNoMoreElement()结束批次。类型映射规则超级表列类型到 SeaTunnel 数据类型的映射由 TDengineTypeMapper.java 实现核心规则如下TDengine 类型SeaTunnel 类型备注BOOL / BITBOOLEANTINYINT / SMALLINT / MEDIUMINT / INT / INTEGER / YEAR含 UNSIGNEDINTINT UNSIGNED / INTEGER UNSIGNED / BIGINTLONGBIGINT UNSIGNEDDECIMAL(20, 0)DECIMAL / DECIMAL UNSIGNEDDECIMAL(38, 18)DECIMAL 会打印可能溢出的告警日志FLOAT / FLOAT UNSIGNEDFLOATUNSIGNED 会打印可能溢出的告警日志DOUBLE / DOUBLE UNSIGNEDDOUBLEUNSIGNED 会打印可能溢出的告警日志CHAR / VARCHAR / TEXT 系列 / JSON / LONGTEXTSTRINGDATELOCAL_DATETIMELOCAL_TIMEDATETIME / TIMESTAMPLOCAL_DATE_TIMEBLOB 系列 / BINARY / VARBINARYBYTESPrimitiveByteArrayTypeGEOMETRY / UNKNOWN不支持抛出UNSUPPORTED_DATA_TYPE异常该映射逻辑有对应的单元测试 TDengineTypeMapperTest.java 验证例如BOOL - BOOLEAN、CHAR - STRING。完整配置示例Source 配置官方文档给出的标准示例TDengine.mdsource { 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.800 result_table_name tdengine_result } }该示例与 e2e 测试场景高度一致。测试 TDengineIT.java 中构造了power.meters超级表包含ts TIMESTAMP, current FLOAT, voltage INT, phase FLOAT, off BOOL五列及location BINARY(64), groupId INT两个标签子表包括d1001~d1004共 8 行测试数据。示例中的时间窗口恰好覆盖这些数据的时间戳从而完整读取 8 条记录。完整任务配置Source Sink以下是一个可运行的完整配置Source 读取 TDengineSink 写入控制台便于验证env { parallelism 4 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.800 result_table_name tdengine_result } } transform { # 可选在此进行过滤、字段映射等转换 } sink { Console { source_result_name tdengine_result } }配置要点result_table_name声明 Source 产出的虚拟表名下游 Sink 通过source_result_name引用每条输出行的第一个字段为subtable_name即子表名如d1001之后依次是超级表定义的列parallelism决定分片分配粒度建议结合子表数量设置让各 Reader 负载均衡若不需要全量数据可通过lower_bound/upper_bound裁剪时间窗口实现增量迁移。运行前提与限制驱动依赖需引入taos-jdbcdriver连接器 pom 中固定为 3.0.3。当 JDBC URL 以jdbc:TAOS-RS://开头时使用 RESTful 驱动对应 taosadapter端口 6041否则使用原生驱动端口 6030必填校验url、database、stable、username、password五项缺失任何一个都会在prepare阶段直接报错批式语义该 Source 为有界BOUNDED批式读取读取完所有分片即结束任务不支持持续消费时间窗口[lower_bound, upper_bound)左闭右开上/下界可单独省略省略一侧则只生成单边过滤条件子表拆分每个子表一个分片子表数量即分片数量上限如果超级表没有子表则不会产生任何数据。总结TDengine Source 连接器为 SeaTunnel 提供了开箱即用的 TDengine 批量读取能力通过url/username/password建立连接按database.stable定位超级表自动发现其下所有子表并按子表拆分并行读取配合lower_bound/upper_bound实现时间窗口迁移。结合 TDengineSourceConfig.java、TDengineSourceSplitEnumerator.java 等源码与 TDengineIT.java 集成测试开发者可以快速上手将 TDengine 中的时序数据迁移到其他存储或分析引擎。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Firebase Source Connector 实战指南从 Firebase Realtime Database 批量同步数据SeaTunnel Firebase Source Connector 实战指南从 Firebase Realtime Database 批量同步数据 本文全数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel GitLab Source Connector 实战指南从 REST API 读取数据与分页同步SeaTunnel GitLab Source Connector 实战指南从 REST API 读取数据与分页同步 本文聚焦 SeaTunnel 的 Git数据集成ETL大数据批处理流处理变更数据捕获OI Wiki 抽象代数入门群、环、域的基本概念与算法竞赛应用OI Wiki 抽象代数入门群、环、域的基本概念与算法竞赛应用 本篇文章以 OI Wiki 数学部分的《代数基础》一章为骨架系统介绍抽象代数中最基础也最常用数据工程大数据批处理流处理上一篇Angular NG05104 根元素未找到Root element was not found错误全解析成因、复现与修复下一篇Dcat Admin 安装与配置终极指南10分钟快速搭建高效后台管理系统创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考