ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

SeaTunnel IoTDBv2 Source 连接器实战指南:从 IoTDB 2.x 树模型/表模型批量读取数据

2026/9/19 3:48:18 拓冰建站 浏览量
SeaTunnel IoTDBv2 Source 连接器实战指南:从 IoTDB 2.x 树模型/表模型批量读取数据 SeaTunnel IoTDBv2 Source 连接器实战指南从 IoTDB 2.x 树模型/表模型批量读取数据【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 仓库中 IoTDBv2 Source 连接器文档 及其对应源码编写讲解如何在 SeaTunnel 作业中通过IoTDBv2连接器从 Apache IoTDB 2.x 读取时序数据覆盖支持引擎、数据类型映射、全部 Source 选项、时间分片并行读取原理并给出树模型与表模型两套可直接运行的 HOCON 配置示例。导读IoTDBv2是 SeaTunnel 面向 Apache IoTDB 2.x 推出的 Source 连接器作业配置中的连接器名称为IoTDBv2。它允许用户直接书写原生 IoTDB SQL 查询语句将查询结果转换成 SeaTunnelRow 后流入下游 Transform 与 Sink同时支持 Spark、Flink 与 SeaTunnel Zeta 三种运行引擎。读完本文你将掌握如何配置IoTDBv2读取树模型与表模型数据、数据在两端如何完成类型映射、如何利用时间列把一次查询拆分成多个分片以充分利用并行度以及每个配置项背后的源码实现原理。支持引擎与主要特性支持引擎根据官方文档说明IoTDBv2Source 支持以下引擎SparkFlinkSeaTunnel Zeta特性清单特性支持情况批处理Batch✅ 支持流处理Streaming✅ 支持精确一次Exactly-once✅ 支持列投影Column Projection✅ 支持IoTDB 通过 SQL 查询天然支持列投影即SELECT子句只选取所需列并行度Parallelism✅ 支持用户自定义分片User-defined Split❌ 暂不支持分片由时间范围自动计算见下文特性定义可参考 Connector V2 特性说明。在源码层面IoTDBv2Source同时实现了SupportParallelism与SupportColumnProjection两个接口对应文档中并行度与列投影两项能力其getBoundedness()返回Boundedness.BOUNDED说明该 Source 本质上是有界读取即每条 SQL 查询执行完毕后读取即结束属于批式数据源。相关实现见 IoTDBv2Source.java。支持的数据源信息数据源支持的版本地址IoTDB2.0 versionlocalhost:6667数据类型映射IoTDBv2Source 将 IoTDB 返回的字段类型转换为 SeaTunnel 数据类型官方映射表如下IoTDB 数据类型SeaTunnel 数据类型BOOLEANBOOLEANINT32TINYINTINT32SMALLINTINT32INTINT64BIGINTFLOATFLOATDOUBLEDOUBLETEXTSTRINGSTRINGSTRINGTIMESTAMPBIGINTTIMESTAMPTIMESTAMPBLOBSTRINGDATEDATE上表看起来一源多映射其转换规则在 DefaultSeaTunnelRowDeserializer.java 中有精确的源码实现核心要点如下INT32 → TINYINT / SMALLINT / INT取决于schema中声明的 SeaTunnel 字段类型。源码中INT32分支会对目标类型做byteValue()TINYINT、shortValue()SMALLINT、intValue()INT三种窄化转换除此之外的声明类型会抛出UNSUPPORTED_DATA_TYPE异常TIMESTAMP → TIMESTAMP / BIGINTIoTDB 时间戳本质是毫秒级long。当 schema 声明为TIMESTAMP时源码将其转换为UTC 时区的LocalDateTimeDate.toInstant().atZone(ZoneOffset.UTC).toLocalDateTime()声明为BIGINT时则直接保留毫秒值DATE直接返回 IoTDB 的DATE对象值BLOB按字符串值读取getStringValue()字段为空时field null对应 SeaTunnel 字段置为null不会导致整行失败。因此schema中声明的类型必须与上表合法组合一致例如 IoTDB 返回INT32时声明为longBIGINT就会在运行时抛出不支持数据类型的异常。Source 选项详解IoTDBv2Source 的全部选项定义在 IoTDBv2SourceOptions.java 中官方文档参数表如下名称类型是否必填默认值描述node_urlsArray是-IoTDB 集群地址格式为[host1:port]或[host1:port,host2:port]usernameString是-IoTDB 用户名passwordString是-IoTDB 用户密码sql_dialectString否treeIoTDB 模型可选值为tree和table。tree表示树模型table表示表模型databaseString否-要查询的数据库名只在表模型中生效sqlString是-要执行的 SQL 查询语句schemaConfig是-数据模式定义详见 Schema 特性fetch_sizeInteger否-单次请求从 IoTDB 获取的行数lower_boundLong否-时间范围下界通过时间列进行数据分片时使用upper_boundLong否-时间范围上界通过时间列进行数据分片时使用num_partitionsInteger否-分区数量通过时间列进行数据分片时使用default_thrift_buffer_sizeInteger否-IoTDB 客户端使用的默认 Thrift 缓冲区大小max_thrift_frame_sizeInteger否-Thrift 最大帧尺寸enable_cache_leaderBoolean否-是否在 IoTDB 客户端启用 Leader 节点缓存common-options否-Source 插件常用参数详见 Source 常用选项连接与执行参数背后的源码实现在 IoTDBv2SourceReader.java 的buildSession()方法中可以看到上述选项是如何驱动 IoTDB 原生 Java Session 的node_urls通过sessionBuilder.nodeUrls(nodes)配置节点列表支持多节点fetch_size通过sessionBuilder.fetchSize(...)控制服务端分批返回的行数直接决定单次网络往返拉取的数据量合理调大可减少 RPC 次数username/password分别设置认证信息default_thrift_buffer_size/max_thrift_frame_size映射到sessionBuilder.thriftDefaultBufferSize(...)与sessionBuilder.thriftMaxFrameSize(...)当单行数据较大或查询返回帧超出默认限制时需要调整enable_cache_leader映射到session.setEnableCacheLeader(...)开启后可减少集群模式下 Leader 节点的寻址开销。每次读取时Reader 对当前分片调用session.executeQueryStatement(split.getQuery())执行 SQL随后遍历SessionDataSet逐行交给反序列化器转换为SeaTunnelRow并output.collect(...)整个连接生命周期open/close由 Reader 管理作业结束后自动关闭 Session。树模型与表模型的内部差异sql_dialect的取值常量定义在 SourceConstants.java 中table与tree。它影响两处行为Reader 选择在 IoTDBv2Source.java 的createReader()中table模型创建IoTDBv2RelationalSourceReader否则创建普通IoTDBv2SourceReader行转换差异树模型下查询结果的第一列是隐含的时间戳RowRecord.getTimestamp()因此DefaultSeaTunnelRowDeserializer.convert()将时间戳写入SeaTunnelRow的第 0 个字段并要求schema字段数 查询列数 1而表模型下convertTableRow()按查询列逐一对应要求schema字段数与查询列数完全一致。这也是两个示例中ts字段位置略有差异的根因。基于时间列的分片并行读取IoTDBv2支持把一条 SQL 按时间列拆分成多个分片Split交由不同 Reader 并行执行从而充分利用env.parallelism配置的并行度。触发条件启用分片读取时需要同时配置lower_bound、upper_bound和num_partitions三个参数只配置num_partitions而缺少上下界或只配置上下界而未配置分区数都不会触发分片。从源码看枚举器在getIotDBSplit()中先判断NUM_PARTITIONS是否配置未配置时直接生成一个分片splitId 为默认值0并使用完整 SQL此时不读取lower_bound/upper_bound。分片算法分片逻辑实现在 IoTDBv2SourceSplitEnumerator.java 的getIotDBSplit()方法中官方文档给出的规则与代码注释完全一致将时间范围分割成 numPartitions 个分区 若 numPartitions 1使用完整的时间范围 若 numPartitions (upper_bound - lower_bound)使用 (upper_bound - lower_bound) 个分区 例lower_bound 1, upper_bound 10, numPartitions 2 sql select * from test where age 0 and age 10 分区结果 split 1: select * from test where (time 1 and time 6) and ( age 0 and age 10 ) split 2: select * from test where (time 6 and time 11) and ( age 0 and age 10 )源码层面的实际实现细节如下SQL 拆分分片前会先把原始 SQL 按where关键字拆成查询主体 条件部分再按align by拆出对齐子句若一条 SQL 包含超过一个where会抛出sql should not contain more than one where异常因此书写分片 SQL 时务必只保留一个where分区数兜底numPartitions通过(end - start) / numPartitions 1计算每段步长size并通过remainder修正边界保证各分区时间区间首尾衔接、不重不漏示例中[1,6)与[6,11)恰好无缝覆盖[1,10]条件拼接每个分片 SQL 查询主体 where (time x and time y) and ( 原始条件 ) align by 子句即在原有查询条件之上叠加时间区间过滤保证切分后语义等价均匀分发分片按splitId排序后通过assignCount % readerCount的轮询方式分配到各并行 ReadergetSplitOwner()使各并行度上的数据量尽量均衡。该行为在 IoTDBv2SourceSplitEnumeratorTest.java 中有shouldBalanceSplitsEvenlyAcrossReaders、shouldContinueRoundRobinAfterRestore、shouldReassignReturnedSplitsToOriginalReader三个测试用例验证覆盖了 4 个 Reader 均分 10 个分片3/3/2/2、checkpoint 恢复后轮询游标延续、失败分片归还给原 Reader 等场景。分片信息splitId与最终查询语句封装在 IoTDBv2SourceSplit.java 中枚举器通过snapshotState()将shouldEnumerate、pendingSplit、assignCount保存到状态配合 IoTDBv2SourceState.java 实现故障恢复后从断点继续分配这是连接器支持精确一次语义的基础。示例一读取 IoTDB 树模型数据以下配置从树模型路径root.test_group.*下按设备align by device读取多列时序数据输出到 Console Sinkenv { parallelism 2 job.mode BATCH } source { IoTDBv2 { node_urls [localhost:6667] username root password root sql SELECT temperature, moisture, c_int, c_bigint, c_float, c_double, c_string, c_boolean FROM root.test_group.* WHERE time 4102329600000 align by device schema { fields { ts timestamp device_name string temperature float moisture bigint c_int int c_bigint bigint c_float float c_double double c_string string c_boolean boolean } } } } sink { Console { } }上游 IoTDB 侧数据格式在 IoTDB CLI 中执行同样的查询返回结果形如IoTDB SELECT temperature, moisture, c_int, c_bigint, c_float, c_double, c_string, c_boolean FROM root.test_group.* WHERE time 4102329600000 align by device; ------------------------------------------------------------------------------------------------------------------------------------- | Time| Device| temperature| moisture| c_int| c_bigint| c_float| c_double| c_string| c_boolean| ------------------------------------------------------------------------------------------------------------------------------------- |2022-09-25T00:00:00.001Z|root.test_group.device_a| 36.1| 100| 1| 21474836470| 1.0f| 1.0d| abc| true| |2022-09-25T00:00:00.001Z|root.test_group.device_b| 36.2| 101| 2| 21474836470| 2.0f| 2.0d| abc| true| |2022-09-25T00:00:00.001Z|root.test_group.device_c| 36.3| 102| 3| 21474836470| 3.0f| 3.0d| abc| true| -------------------------------------------------------------------------------------------------------------------------------------读取到 SeaTunnelRow 后的格式树模型下schema首字段ts接收 IoTDB 的行时间戳毫秒longdevice_name对应align by device产生的设备列其余字段按查询列顺序对应tsdevice_nametemperaturemoisturec_intc_bigintc_floatc_doublec_stringc_boolean1664035200001root.test_group.device_a36.11001214748364701.0f1.0dabctrue1664035200001root.test_group.device_b36.21012214748364702.0f2.0dabctrue1664035200001root.test_group.device_c36.31023214748364703.0f3.0dabctrue注意时间戳由2022-09-25T00:00:00.001Z变为毫秒值1664035200001这正是前文TIMESTAMP → BIGINT映射的体现schema中ts timestamp声明为timestamp类型时源码会按 UTC 时区转换为LocalDateTime。示例二读取 IoTDB 表模型数据以下配置通过sql_dialect table读取表模型数据并用database指定目标数据库env { parallelism 2 job.mode BATCH } source { IoTDBv2 { node_urls [localhost:6667] username root password root sql_dialect table database test_database sql SELECT time, sn, type, bidprice, bidsize, domain, buyno, askprice FROM test_table schema { fields { ts timestamp sn string type string bidprice int bidsize double domain boolean buyno bigint askprice string } } } } sink { Console { } }提示若查询语句中已明确写出数据库如FROM test_database.test_table则无需再配置database参数。上游 IoTDB 侧数据格式IoTDB SELECT time, sn, type, bidprice, bidsize, domain, buyno, askprice FROM test_table --------------------------------------------------------------------------------------- | time| sn|type|bidprice| bidsize|domain|buyno| askprice| --------------------------------------------------------------------------------------- |2025-07-30T17:52:34.85108:00|0700HK| L1| 9|10.323907796459721| true| 10|-1064754527| |2025-07-30T17:52:34.95108:00|0700HK| L1| 10| 9.844574317657585| false| 9|-1088662576| |2025-07-30T17:52:35.05108:00|0700HK| L1| 9| 9.272974132434069| true| 9| 402003616| ---------------------------------------------------------------------------------------读取到 SeaTunnelRow 后的格式表模型下time列是普通查询列schema按查询列一一对应字段数与查询列数一致tssntypebidpricebidsizedomainbuynoaskprice2025-07-30T17:52:34.8510700HKL1910.323907796459721true10-10647545272025-07-30T17:52:34.9510700HKL1109.844574317657585false9-10886625762025-07-30T17:52:35.0510700HKL199.272974132434069true9402003616由于schema中ts timestamp声明为timestamp类型时区08:00的时间在转换时被统一归一化为 UTC 表示2025-07-30T17:52:34.851。常见问题与最佳实践sql中避免多个where若打算使用时间分片原始 SQL 只允许出现一个where否则枚举器会抛出sql should not contain more than one where异常如需额外过滤条件请将其合并在同一个where中分片时会自动以and拼接schema字段数与查询列严格匹配树模型下schema字段数 查询列数 1首字段接收时间戳表模型下schema字段数 查询列数不一致会在反序列化时抛出Illegal SeaTunnelRowType异常合理设置fetch_size与 Thrift 参数大批量或大字段如 BLOB查询时适当调大fetch_size可减少往返次数必要时同步调大default_thrift_buffer_size/max_thrift_frame_size以避免帧超限分片与并行度配合分片数决定可被并行执行的任务数建议分片数不小于env.parallelism让每个 Reader 都能分配到分片num_partitions应结合时间范围与数据量设置避免分区过碎造成额外开销集群部署可开启enable_cache_leader在 IoTDB 集群模式下开启 Leader 缓存可减少寻址开销单机模式下无影响。延伸阅读连接器整体架构与特性体系Connector V2 特性说明schema配置细则Schema 特性Source 插件公共参数Source 常用选项连接器变更记录IoTDB 连接器 Changelog源码与测试选项定义见 IoTDBv2SourceOptions.java分片实现见 IoTDBv2SourceSplitEnumerator.java分片均衡分配测试见 IoTDBv2SourceSplitEnumeratorTest.java【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考