SeaTunnel TDengine Source Connector 实战指南:配置、原理与数据同步

发布时间:2026/9/29 9:24:31

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),仅供参考
延伸阅读

更多相关文章

2026/9/29 9:24:31

Paperclip范式:轻量级AI Agent协同开发实践指南

1. “Paperclip”不是回形针:它是一套AI Agent协同开发范式你搜“paperclip”,第一反应是办公桌抽屉里那枚银色小金属?别急——在2024年中后期的开发者社区里,paperclip 已悄然成为一类轻量级、可组合、面向任务流的AI Agent开发范…

2026/9/29 9:24:31

从零构建大模型:AI工程全链路实战指南

1. "从零开始"到底意味着什么:先想清楚再动手这两年"AI engineering"几乎成了最热的岗位关键词,而GitHub上诸如build a reasoning model from scratch、build a large language model from scratch这类项目动辄上万star。很多人一上…

2026/9/29 10:29:37

《Python 大型项目重构实战指南:从混乱到秩序的渐进式策略》 - 实践

《大型项目重构实战指南: 从混乱到秩序的渐进式策略》这本书, 它的第一个大项就是引言部分, 引言部分的主要内容是在回答为什么要进行重构。在软件开发这个领域, 根本不可能存在哪一个项目是一直永远都能够保持那种“干净整洁”的状态的。后面会因为业务上的那些需要被非常快速…

2026/9/29 10:29:37

实例分割避坑指南:从原理到工程实践的完整心得

1. 为什么我劝你先想清楚再入坑实例分割聊实例分割之前,先讲个真实感受。我接触这个方向大概一年半,从最开始对着 Mask R-CNN 的论文看得一头雾水,到后来能独立调通 SOLO、跑通 YOLO 系列的实例分割分支,中间踩过的坑比代码行数还…

2026/9/29 10:29:37

WorkBuddy+腾讯乐享:企业知识Agent化工作流重构实践

1. 这不是又一个“知识库接入教程”,而是工作流重构的临界点我第一次在腾讯乐享后台看到那个“WorkBuddy 接入”按钮时,下意识点了右上角的叉——以为又是某个需要填三页表单、配五层权限、等三天审批的“企业级集成”。直到上周,市场部同事甩…

2026/9/29 10:29:37

Agent从能说到能干:Skill技能系统设计与避坑实战

做了这么久的 Agent 相关项目,我收到最多的反馈其实是同一句话:“它聊得头头是道,我让它干活怎么就这么费劲?” 这不是模型不行,而是我们一直在用“聊天”的思路去要求“干活”。这个系列写到第八篇,前面聊…

2026/9/29 10:29:37

Dify完全指南:安装部署、知识库与Agent工作流实战

最近很多人在问 Dify 是什么、Dify 怎么装、Dify 能做什么。我接触 Dify 也有一年多了,从 0.6 版本一路用到社区版 1.10,中间踩过不少坑,也看着它从一个大模型管理面板慢慢长成现在这个集知识库、工作流、Agent、可观测性于一体的小型应用平台…

2026/9/29 10:24:37

AI视频多轨道怎么同步预览:统一时钟、找锚点并排查偏移

AI视频多轨道同步预览,首先要让画面、声音、字幕、特效和转场在同一条时间线的时间基准下参与合成,再用可识别的动作或声音事件检查彼此关系。预览窗口里有画面在播放,只说明部分内容可播放,不证明全部轨道已经同步,也…

2026/9/28 3:03:23

东莞市品牌网站建设报价常见报错与解决

东莞品牌网站建设报价单背后:一份保姆级建站教程避坑实录 网站做好了没人访问,这大概是很多老板最头疼的事。花了大几万做的品牌站,上线后流量惨淡,比路边摊还冷清。别急着骂外包公司,很多“东莞品牌网站建设报价”里藏着不少猫腻,比如用模板站冒充定制…

2026/9/28 6:05:15

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解 【免费下载链接】spirula-studio Cross-vendor 3D Gaussian Splatting trainer - video to splat to mesh, Vulkan or CUDA. 项目地址: https://gitcode.com/GitHub_Trending/sp/spirula-studio Sp…

2026/9/29 7:00:49

SEO怎么推广速查手册新手避坑实战指南

SEO怎么推广速查手册新手避坑实战指南 模板网站太丑不够用?别急着加滤镜,那是治标不治本。很多老板盯着后台流量掉得眼红,却还在纠结首页Banner的圆角是不是3像素。这就像穿着西装去挖土,姿势不对,努力白费。我整理这份 速查手册…

2026/9/29 0:04:04

AI Evals实战指南:从零搭建LLM应用评估体系与CI/CD集成

1. 为什么AI Evals值得你花时间搞明白做LLM应用的人,迟早会撞上同一堵墙:模型输出飘忽不定,今天答得好好的,明天换个问法就胡说八道。你改了一版提示词,感觉好像好了点,但到底好了多少?说不清。…

2026/9/29 0:04:04

Java采购管理系统实战:从数据库设计到事务一致性

简介:这是一套面向Java Web初学者与课程设计者的采购管理系统完整源码,采用JSP技术搭建,配合MySQL数据库,用于解决企业采购信息的管理问题,适合作为毕业设计、课程大作业或进销存类项目的参考模板。系统实现了用户登录…

2026/9/29 3:53:39

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/29 9:46:12

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/29 6:36:14

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
☎咨询二维码 ☎ ↑