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 支持以下引擎:
- Spark
- Flink
- SeaTunnel 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。
支持的数据源信息
| 数据源 | 支持的版本 | 地址 |
|---|---|---|
| IoTDB | 2.0 <= version | localhost:6667 |
数据类型映射
IoTDBv2Source 将 IoTDB 返回的字段类型转换为 SeaTunnel 数据类型,官方映射表如下:
| IoTDB 数据类型 | SeaTunnel 数据类型 |
|---|---|
| BOOLEAN | BOOLEAN |
| INT32 | TINYINT |
| INT32 | SMALLINT |
| INT32 | INT |
| INT64 | BIGINT |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| TEXT | STRING |
| STRING | STRING |
| TIMESTAMP | BIGINT |
| TIMESTAMP | TIMESTAMP |
| BLOB | STRING |
| DATE | DATE |
上表看起来"一源多映射",其转换规则在 DefaultSeaTunnelRowDeserializer.java 中有精确的源码实现,核心要点如下:
- INT32 → TINYINT / SMALLINT / INT:取决于
schema中声明的 SeaTunnel 字段类型。源码中INT32分支会对目标类型做byteValue()(TINYINT)、shortValue()(SMALLINT)、intValue()(INT)三种窄化转换,除此之外的声明类型会抛出UNSUPPORTED_DATA_TYPE异常; - TIMESTAMP → TIMESTAMP / BIGINT:IoTDB 时间戳本质是毫秒级
long。当 schema 声明为TIMESTAMP时,源码将其转换为UTC 时区的LocalDateTime(Date.toInstant().atZone(ZoneOffset.UTC).toLocalDateTime());声明为BIGINT时则直接保留毫秒值; - DATE:直接返回 IoTDB 的
DATE对象值; - BLOB:按字符串值读取(
getStringValue()); - 字段为空时(
field == null)对应 SeaTunnel 字段置为null,不会导致整行失败。
因此,schema中声明的类型必须与上表合法组合一致,例如 IoTDB 返回INT32时声明为long(BIGINT)就会在运行时抛出不支持数据类型的异常。
Source 选项详解
IoTDBv2Source 的全部选项定义在 IoTDBv2SourceOptions.java 中,官方文档参数表如下:
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| node_urls | Array | 是 | - | IoTDB 集群地址,格式为["host1:port"]或["host1:port","host2:port"] |
| username | String | 是 | - | IoTDB 用户名 |
| password | String | 是 | - | IoTDB 用户密码 |
| sql_dialect | String | 否 | tree | IoTDB 模型,可选值为tree和table。tree表示树模型,table表示表模型 |
| database | String | 否 | - | 要查询的数据库名,只在表模型中生效 |
| sql | String | 是 | - | 要执行的 SQL 查询语句 |
| schema | Config | 是 | - | 数据模式定义,详见 Schema 特性 |
| fetch_size | Integer | 否 | - | 单次请求从 IoTDB 获取的行数 |
| lower_bound | Long | 否 | - | 时间范围下界(通过时间列进行数据分片时使用) |
| upper_bound | Long | 否 | - | 时间范围上界(通过时间列进行数据分片时使用) |
| num_partitions | Integer | 否 | - | 分区数量(通过时间列进行数据分片时使用) |
| default_thrift_buffer_size | Integer | 否 | - | IoTDB 客户端使用的默认 Thrift 缓冲区大小 |
| max_thrift_frame_size | Integer | 否 | - | Thrift 最大帧尺寸 |
| enable_cache_leader | Boolean | 否 | - | 是否在 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的轮询方式分配到各并行 Reader(getSplitOwner()),使各并行度上的数据量尽量均衡。该行为在 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 Sink:
env { 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 的行时间戳(毫秒long),device_name对应align by device产生的设备列,其余字段按查询列顺序对应:
| ts | device_name | temperature | moisture | c_int | c_bigint | c_float | c_double | c_string | c_boolean |
|---|---|---|---|---|---|---|---|---|---|
| 1664035200001 | root.test_group.device_a | 36.1 | 100 | 1 | 21474836470 | 1.0f | 1.0d | abc | true |
| 1664035200001 | root.test_group.device_b | 36.2 | 101 | 2 | 21474836470 | 2.0f | 2.0d | abc | true |
| 1664035200001 | root.test_group.device_c | 36.3 | 102 | 3 | 21474836470 | 3.0f | 3.0d | abc | true |
注意时间戳由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.851+08:00|0700HK| L1| 9|10.323907796459721| true| 10|-1064754527| |2025-07-30T17:52:34.951+08:00|0700HK| L1| 10| 9.844574317657585| false| 9|-1088662576| |2025-07-30T17:52:35.051+08:00|0700HK| L1| 9| 9.272974132434069| true| 9| 402003616| +-----------------------------+------+----+--------+------------------+------+-----+-----------+读取到 SeaTunnelRow 后的格式
表模型下time列是普通查询列,schema按查询列一一对应(字段数与查询列数一致):
| ts | sn | type | bidprice | bidsize | domain | buyno | askprice |
|---|---|---|---|---|---|---|---|
| 2025-07-30T17:52:34.851 | 0700HK | L1 | 9 | 10.323907796459721 | true | 10 | -1064754527 |
| 2025-07-30T17:52:34.951 | 0700HK | L1 | 10 | 9.844574317657585 | false | 9 | -1088662576 |
| 2025-07-30T17:52:35.051 | 0700HK | L1 | 9 | 9.272974132434069 | true | 9 | 402003616 |
由于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),仅供参考