SeaTunnel IoTDB 连接器实战指南:从时序数据读取到批量写入

发布时间:2026/9/16 10:42:21
SeaTunnel IoTDB 连接器实战指南:从时序数据读取到批量写入 SeaTunnel IoTDB 连接器实战指南从时序数据读取到批量写入【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelApache SeaTunnel本仓库是一个多模态、高性能、分布式的海量数据集成工具。其中 IoTDB 连接器负责打通 Apache IoTDB 时序数据库与 SeaTunnel 数据管道作为 Source 从 IoTDB 执行有界 SQL 查询读取时序数据作为 Sink 将 SeaTunnelRow 批量写入 IoTDB。本文基于仓库中的 IoTDB Source 文档 与 IoTDB Sink 文档其变更历史统一维护在 connector-iotdb changelog结合连接器源码与测试系统讲解连接器的能力边界、配置项、数据类型映射、时间分区切分原理、多表读取以及完整的读写实战配置帮助读者在 SeaTunnel 中快速落地 IoTDB 数据的集成方案。支持引擎与能力边界该连接器同时支持 Spark、Flink 与 SeaTunnel Zeta 三种执行引擎。Source 能力特性根据 IoTDB Source 文档 中的 Key features 列表Source 具备以下能力✅ 批处理batch❌ 流式stream❌ 精确一次exactly-once✅ 列投影column projectionIoTDB 本身支持通过 SQL 查询完成列投影✅ 并行度parallelism✅ 多表multiple table❌ 用户自定义分片support user-defined split❌ CDC关键限制IoTDB Source 每个任务只会执行一条有界bounded的 SQL 查询适用于批作业或有界时间窗读取。它不会在 IoTDB 的变更日志上维护流式游标因此无法在流式模式下持续尾随tail新写入的数据。这一点在源码中也能得到印证——IoTDBSource.java 中getBoundedness()直接返回Boundedness.BOUNDED。Sink 能力特性根据 IoTDB Sink 文档 的 Key Features✅ 精确一次exactly-once通过幂等写入实现。如果多条数据具有相同的keydevice和timestamp后写入的值会覆盖先前的值。❌ CDC写入语义说明Sink 通过调用 IoTDB 的 insert RPC 写行。当某一行携带非唯一的(device, timestamp)组合时写入被视为 upsert——最新值覆盖旧值因此上游重复投递不会产生幻影行。行类型的UPDATE/DELETE不会被解释为 CDC 操作所有行都以 insert 方式写入。这一行为与连接器源码中的实现一致见 IoTDBSinkClient.java。支持的数据源版本与数据类型映射Source 与 Sink 都支持 IoTDB 版本0.13.0 version 1.3.X默认连接地址示例为localhost:6667IoTDB 默认 Thrift RPC 端口。数据类型映射IoTDB 数据类型SeaTunnel 数据类型BOOLEANBOOLEANINT32TINYINTINT32SMALLINTINT32INTINT64BIGINTFLOATFLOATDOUBLEDOUBLETEXTSTRINGSTRINGSTRINGtime 列BIGINTtime 列TIMESTAMP从映射表可以看到IoTDB 的 INT32 可以按精度需求映射为 SeaTunnel 的 TINYINT、SMALLINT 或 INT时间列既可以用bigintepoch 毫秒表示也可以用timestamp表示。Source 选项详解使用 Source 时可以选择根级sqlschema的单表配置或**tables_configs的多表独立配置**原有的单表配置方式保持兼容。名称类型必填默认值描述node_urlsstring是-IoTDB 集群地址格式为host1:port或host1:port,host2:portusernamestring是-IoTDB 用户名passwordstring是-IoTDB 密码sqlstring条件必填-SQL 查询。未配置tables_configs时与schema搭配必填tables_configsarray否-表配置列表每个条目包含sql和schema且schema.table必须唯一且非空schemaconfig条件必填-与根级sql搭配必填使用tables_configs时在每个条目内配置fetch_sizeint否-单次请求从 IoTDB 拉取的行数lower_boundlong否-SeaTunnel 按时间切分查询时的下界upper_boundlong否-SeaTunnel 按时间切分查询时的上界num_partitionsint否-时间范围分区数需与lower_bound、upper_bound配合使用thrift_default_buffer_sizeint否-IoTDB 客户端的初始 Thrift 缓冲区大小thrift_max_frame_sizeint否-IoTDB 客户端的最大 Thrift 帧大小enable_cache_leaderboolean否-IoTDB 客户端是否缓存 leader 节点versionstring否-客户端使用的 SQL 语义版本可选V_0_12和V_0_13common-options-否-Source 插件通用参数参见 Source Common Optionsschema.fields中的第一个字段必须描述 IoTDB 时间列可以是bigintepoch 毫秒或timestampSeaTunnel 时间戳值。从源码看 Source 会话构建在 IoTDBSourceReader.java 的buildSession()方法中连接器使用Session.Builder构建 IoTDB 会话将node_urls按分隔符拆分成节点列表然后把fetch_size、username、password、thrift_default_buffer_size、thrift_max_frame_size、enable_cache_leader、version逐一透传给 Session Builder。所有参数的定义集中在 IoTDBSourceOptions.java。读取过程read()方法通过session.executeQueryStatement(split.getQuery())执行分片后的 SQL遍历SessionDataSet用DefaultSeaTunnelRowDeserializer将每一行RowRecord反序列化为SeaTunnelRow并为多表场景设置tableId后交给下游 Collector。Sink 选项详解名称类型必填默认值描述node_urlsArray是-IoTDB 集群地址格式为[host1:port]或[host1:port,host2:port]usernameString是-IoTDB 用户名passwordString是-IoTDB 密码key_deviceString是-指定 SeaTunnelRow 中承载 IoTDB deviceId 的字段名key_timestampString否processing time指定 SeaTunnelRow 中承载 IoTDB timestamp 的字段名未指定时使用处理时间processing timekey_measurement_fieldsArray否排除 device 与 timestamp 字段外的全部字段指定 SeaTunnelRow 中 IoTDB measurement 列表的字段名storage_groupString否-指定设备存储组路径前缀。例deviceId ${storage_group} . ${key_device}batch_sizeInteger否1024批量写入当缓冲行数达到batch_size时刷入 IoTDBmax_retriesInteger否-失败刷写重试次数retry_backoff_multiplier_msInteger否-生成下一次退避延迟的乘数max_retry_backoff_msInteger否-重试请求 IoTDB 前等待的最大时间default_thrift_buffer_sizeInteger否-IoTDB 客户端的 Thrift 初始缓冲区大小max_thrift_frame_sizeInteger否-IoTDB 客户端的 Thrift 最大帧大小zone_idstring否-IoTDB 客户端使用的 java.time.ZoneIdenable_rpc_compressionBoolean否-是否启用 IoTDB 客户端 RPC 压缩connection_timeout_in_msInteger否-连接 IoTDB 的最大等待时间毫秒common-options-否-Sink 插件通用参数参见 Sink Common Options从 IoTDBSinkOptions.java 的源码可以看到batch_size的默认值确实是1024DEFAULT_BATCH_SIZE且该选项声明了默认值其余选项均为noDefaultValue()。Sink 写入规则key_device必须指定 SeaTunnel 中包含 IoTDB 设备路径的字段。storage_group是字符串前缀设置后最终设备路径由storage_group与key_device字段的值拼接而成。key_timestamp可以是STRING、BIGINT或TIMESTAMP类型字段未配置时使用当前处理时间。若未配置key_measurement_fields除key_device和key_timestamp外的所有字段都会被写为 measurement。Sink 支持的 measurement 字段类型包括STRING、BOOLEAN、TINYINT、SMALLINT、INT、BIGINT、FLOAT、DOUBLE。Sink 客户端工作机制从源码看IoTDBSinkClient.java 的write()方法先把记录加入batchList当缓冲行数达到batch_size时触发flush()批量写入 IoTDBclose()时会做最后一次 flush 并关闭会话。这与文档中「缓冲行数达到batch_size时刷入」的描述一致。文档还指出对于长时间运行的流作业连接器会在缓冲区满或checkpoint 完成时刷写缓冲行。时间范围分区Time-Range Partitioninglower_bound、upper_bound、num_partitions三个参数让 SeaTunnel 将单条有界查询拆分为多个子查询每个子查询对应一个时间分区。这些分区会按照配置的parallelism分发到各个子任务每个子任务负责互不重叠的时间片。分区规则若num_partitions 1整个范围[lower_bound, upper_bound)作为一个分区使用。否则将范围均分为num_partitions份如果upper_bound - lower_bound num_partitions则退化为upper_bound - lower_bound个分区。例如lower_bound 1、upper_bound 10、num_partitions 2针对 SQLselect * 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);从 IoTDBSourceSplitEnumerator.java 的getIotDBSplit()方法可以印证该实现先按align by与where关键字拆分原始 SQL一条 SQL 至多含一个 where再按(end - start) / numPartitions 1计算每个分区大小生成where (time X and time Y)的时间谓词并与原条件、align by子句重组为分片查询。分区以轮询方式assignCount % numReaders分配给各 reader实现并行读取。适用场景当时间列是天然分区键、且底层数据集跨越很大时间范围时使用时间分区最合适。对于非时间维度的并行化需求直接提高parallelism并依赖 IoTDB 侧的分区能力即可。多表读取Multiple Tables使用tables_configs可以在一个 Source 中配置多张表。要点如下所有表共享 Source 级的连接与客户端选项包括node_urls、用户名密码、fetch_size、version等。这些选项只能在 Source 级配置。每个tables_configs条目定义自己的 SQL 查询、schema 以及可选的lower_bound、upper_bound、num_partitions。根级 SQL、schema 与时间分区选项不能与tables_configs混用。schema.table标识输出表用于下游路由它不会改变 SQL 中的 IoTDB 路径。每个 schema 必须与其查询的列按顺序匹配且以时间列开头。不同条目可以有不同的字段名与类型。配置示例source { IoTDB { node_urls localhost:6667 username root password root tables_configs [ { sql SELECT temperature FROM root.weather.device_a lower_bound 1 upper_bound 100 num_partitions 4 schema { table weather fields {ts bigint, temperature float} } }, { sql SELECT enabled FROM root.status.device_b schema { table status fields {ts bigint, enabled boolean} } } ] } }多表时间分区约束对于带时间分区的表需要同时提供三个分区选项、正数的num_partitions且lower_bound upper_bound。范围在两个边界处均为闭区间被切分为互不重叠的分区每个时间戳最多产生一个分区。不支持upper_bound Long.MAX_VALUE以及闭区间大小超过Long.MAX_VALUE的范围。不设置分区选项时该表的 SQL 会作为一个 split 原样执行根级分区行为保持不变。逐表时间分区只支持简单的 SELECT 投影可带 WHERE 与 ALIGN BY 子句。以下 SQL 结构在此模式下会被拒绝因为追加时间谓词无法安全保持其语义SELECT 列表中的函数、子查询、带引号的表达式/标识符、SQL 注释、分号、GROUP BY、ORDER BY、LIMIT/OFFSET、SLIMIT/SOFFSET、FILL 与 INTO。遇到这类查询时不设置分区选项即可原样执行。多表与检查点多表模式仍然是有界 Source不是 CDC。检查点会保留每个待处理分片的表标识。旧的单表检查点在原单表配置下仍然可读但如果要把一个运行中的作业从单表配置切换为tables_configs需要启动新作业而不能恢复旧检查点。其他要点当 SQL 使用align by device时第二个字段通常描述 IoTDB 设备名其余字段必须与 SQL 查询返回的 measurement 顺序一致也可以将时间列用作 SQL 查询中的分区键。这些行为在多表测试 IoTDBMultiTableSourceTest.java 中有所覆盖。Source 实战示例示例一基础读取含时间分区以下配置从 IoTDB 读取多个设备的时序数据并开启时间分区env { parallelism 2 job.mode BATCH } source { IoTDB { 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 lower_bound 1 upper_bound 4102329600000 num_partitions 10 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 { } }lower_bound、upper_bound、num_partitions均为可选参数。当查询覆盖很大时间范围、希望 SeaTunnel 把读取拆成多个时间分区并行执行时它们非常有用。IoTDB 中align by device查询返回的数据格式如下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 后的数据格式时间列被解析为 epoch 毫秒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示例二IoTDB 到 IoTDB结合 Replace 转换下面的示例从一个 IoTDB 路径读取数据替换设备前缀后写入另一个 IoTDB 路径演示了 Source、Transform、Sink 三者的联动env { parallelism 2 job.mode BATCH } source { IoTDB { plugin_output iotdb_rows node_urls localhost:6667 username root password root sql SELECT c_string, c_boolean, c_tinyint, c_smallint, c_int, c_bigint, c_float, c_double FROM root.source_group.* WHERE time 4102329600000 align by device lower_bound 1 upper_bound 4102329600000 num_partitions 10 schema { fields { ts timestamp device_name string c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double } } } } transform { Replace { plugin_input iotdb_rows plugin_output sink_rows replace_field device_name pattern root.source_group replacement root.sink_group is_regex false replace_first true } } sink { IoTDB { plugin_input sink_rows node_urls [localhost:6667] username root password root key_device device_name key_timestamp ts key_measurement_fields [c_string, c_boolean, c_tinyint, c_smallint, c_int, c_bigint, c_float, c_double] batch_size 1 } }Sink 实战示例Sink 部分先看一个典型的 FakeSource 造数配置它构造了包含设备、测量值与事件时间的 SeaTunnelRowenv { parallelism 2 job.mode BATCH } source { FakeSource { row.num 16 bigint.template [1664035200001] schema { fields { device_name string temperature float moisture int event_ts bigint c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double } } } }上游 SeaTunnelRow 的数据格式如下device_nametemperaturemoistureevent_tsc_stringc_booleanc_tinyintc_smallintc_intc_bigintc_floatc_doubleroot.test_group.device_a36.11001664035200001abc1true11121474836481.01.0root.test_group.device_b36.21011664035200001abc2false22221474836492.02.0root.test_group.device_c36.31021664035200001abc3false33321474836493.03.0Case 1仅使用必填选项不指定key_timestamp使用当前处理时间作为时间戳measurement 字段包含除key_device外的所有字段sink { IoTDB { node_urls [localhost:6667] username root password root key_device device_name # 指定 deviceId 使用 device_name 字段 } }IoTDB 输出格式IoTDB SELECT * FROM root.test_group.* align by device; ----------------------------------------------------------------------------------------------------------------------------------------------------------------- | Time| Device| temperature| moisture| event_ts| c_string| c_boolean| c_tinyint| c_smallint| c_int| c_bigint| c_float| c_double| ----------------------------------------------------------------------------------------------------------------------------------------------------------------- |2023-09-01T00:00:00.001Z|root.test_group.device_a| 36.1| 100| 1664035200001| abc1| true| 1| 1| 1| 2147483648| 1.0| 1.0| |2023-09-01T00:00:00.001Z|root.test_group.device_b| 36.2| 101| 1664035200001| abc2| false| 2| 2| 2| 2147483649| 2.0| 2.0| |2023-09-01T00:00:00.001Z|root.test_group.device_c| 36.3| 102| 1664035200001| abc2| false| 3| 3| 3| 2147483649| 3.0| 3.0| -----------------------------------------------------------------------------------------------------------------------------------------------------------------注意这里 Time 列显示为写入时刻的处理时间2023-09-01而不是事件时间。Case 2使用上游事件时间通过key_timestamp使用事件时间measurement 字段包含除key_device与key_timestamp外的所有字段sink { IoTDB { node_urls [localhost:6667] username root password root key_device device_name # 指定 deviceId 使用 device_name 字段 key_timestamp event_ts # 指定 timestamp 使用 event_ts 字段 } }此时 IoTDB 中 Time 列即为上游event_ts对应的时间2022-09-25与 Case 1 的处理时间形成对比IoTDB SELECT * FROM root.test_group.* align by device; ----------------------------------------------------------------------------------------------------------------------------------------------------------------- | Time| Device| temperature| moisture| event_ts| c_string| c_boolean| c_tinyint| c_smallint| c_int| c_bigint| c_float| c_double| ----------------------------------------------------------------------------------------------------------------------------------------------------------------- |2022-09-25T00:00:00.001Z|root.test_group.device_a| 36.1| 100| 1664035200001| abc1| true| 1| 1| 1| 2147483648| 1.0| 1.0| |2022-09-25T00:00:00.001Z|root.test_group.device_b| 36.2| 101| 1664035200001| abc2| false| 2| 2| 2| 2147483649| 2.0| 2.0| |2022-09-25T00:00:00.001Z|root.test_group.device_c| 36.3| 102| 1664035200001| abc2| false| 3| 3| 3| 2147483649| 3.0| 3.0| -----------------------------------------------------------------------------------------------------------------------------------------------------------------Case 3限定 measurement 字段同时使用key_timestamp与key_measurement_fields只把指定字段写入 IoTDBsink { IoTDB { node_urls [localhost:6667] username root password root key_device device_name key_timestamp event_ts key_measurement_fields [temperature, moisture] } }IoTDB 输出只保留temperature与moisture两个测量列IoTDB SELECT * FROM root.test_group.* align by device; ------------------------------------------------------------------------- | Time| Device| temperature| moisture| ------------------------------------------------------------------------- |2022-09-25T00:00:00.001Z|root.test_group.device_a| 36.1| 100| |2022-09-25T00:00:00.001Z|root.test_group.device_b| 36.2| 101| |2022-09-25T00:00:00.001Z|root.test_group.device_c| 36.3| 102| -------------------------------------------------------------------------Case 4流式写入与显式批量刷写对于长时间运行的流作业调大batch_size可以降低逐行 RPC 的开销。连接器在缓冲区填满batch_size或 checkpoint 完成时刷写缓冲行。配置max_retries与max_retry_backoff_ms可以让作业对瞬时 RPC 故障更具韧性env { parallelism 2 job.mode STREAMING checkpoint.interval 10000 } sink { IoTDB { node_urls [localhost:6667, localhost:6668] username root password root key_device device_name key_timestamp event_ts batch_size 2048 max_retries 3 retry_backoff_multiplier_ms 100 max_retry_backoff_ms 5000 } }node_urls支持配置多个 IoTDB 节点。Sink 会为每个任务挑选一个活跃写入节点当活跃节点故障时自动切换到其他节点。连接器演进与维护线索IoTDB 连接器的变更历史统一记录在 connector-iotdb changelog 中从中可以看到该连接器的主要演进脉络2.2.0-beta新增 IoTDB Sink#2407与 IoTDB Source#2431连接器正式诞生。2.3.0 / 2.3.0-beta统一 IoTDB Source 与 Sink 的异常处理#3557、为 IoTDB 暴露可配置选项#3387、增加 Sink 参数校验#3412、修复 Sink NPE#3080等。2.3.1统一 schema 参数并更新 IoTDB Source#3896、增加并行度与列投影接口#3829、重构 schema 解析#4157。2.3.3 / 2.3.4修复 shade 后 Preconditions 引用问题#5284、移除 IoTDB Sink 中的 scheduler#5270、支持 schema 中配置 column/primaryKey/constraintKey#5564、引入新的错误定义规则#5793。2.3.9 / 2.3.10增加 shade 检查规则#8136、重构 connector common options#8634、改进 iotdb options#8965。2.3.11为 source/sink 状态类缺失serialVersionUID增加检查脚本#9118。这些变更表明连接器在稳定性异常处理、参数校验、可配置性选项暴露与代码质量shade、序列化检查三个方向持续演进当前文档所述的多表tables_configs配置即属于较新的能力。总结SeaTunnel 的 IoTDB 连接器为时序数据集成提供了完整闭环Source 侧通过有界 SQL 查询读取 IoTDB 数据支持列投影、并行度、时间范围分区与多表读取Sink 侧通过幂等写入实现 exactly-once 语义支持设备/时间戳/measurement 的灵活映射、批量刷写与多节点故障切换。使用前需明确其能力边界——Source 是批式有界读取而非流式 CDC而 Sink 的(device, timestamp)唯一性保证正好适用于时序场景的乱序重放与去重。如需进一步了解 schema 定义细节可参考 Schema Feature 与 Source/Sink Common Options。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询