SeaTunnel IoTDBv2 Sink 连接器实战:树模型与表模型双模式写入的原理、参数详解与配置示例 📅 发布时间:2026/9/19 2:25:50 👁 浏览次数: SeaTunnel IoTDBv2 Sink 连接器实战树模型与表模型双模式写入的原理、参数详解与配置示例【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文以 docs/zh/connectors/sink/IoTDBv2.md 为核心脉络结合 Apache SeaTunnel 仓库中connector-iotdb-v2模块的源码实现系统讲解 IoTDBv2 数据接收器的整体能力、字段映射规则、全部 Sink 选项、树模型tree与表模型table两种写入模式的完整配置示例以及底层批量写入与重试机制。读完本文你将能够在 SeaTunnel 作业中正确配置 IoTDBv2 Sink将 SeaTunnelRow 数据写入 IoTDB 2.x并能根据数据模型灵活选用key_device、key_timestamp、key_measurement_fields、key_tag_fields、key_attribute_fields等关键参数。连接器概述与支持引擎IoTDBv2 是 SeaTunnel 用于将数据写入 Apache IoTDB 2.x 的数据接收器作业配置中的连接器名称为IoTDBv2大小写敏感。它支持以下运行引擎SparkFlinkSeaTunnel ZetaIoTDBv2 Sink 的主要特性如下✅ 精确一次Exactly-OnceIoTDB 通过幂等写支持精确一次。如果两条数据使用相同的key和timestamp新数据将覆盖旧数据从而在重复提交时不会产生重复记录。❌ 定时刷新Periodically Flush暂不支持。支持的数据源信息数据源支持的版本默认地址IoTDB2.0 versionlocalhost:6667从源码看该连接器模块位于seatunnel-connectors-v2/connector-iotdb-v2同时包含 Sinksink/包与 Sourcesource/包两套实现本文聚焦 Sink 方向。SeaTunnel 与 IoTDB 数据类型映射IoTDBv2 Sink 内置了从 SeaTunnel 类型到 IoTDB 类型的转换逻辑映射关系如下SeaTunnel 数据类型IoTDB 数据类型BOOLEANBOOLEANTINYINTINT32SMALLINTINT32INTINT32BIGINTINT64FLOATFLOATDOUBLEDOUBLESTRINGSTRINGTIMESTAMPTIMESTAMPDATEDATE该映射表可以在 DefaultSeaTunnelRowSerializer.java 的convert(SeaTunnelDataType)方法中得到源码印证TINYINT、SMALLINT、INT统一转换为TSDataType.INT32BIGINT→TSDataType.INT64FLOAT→TSDataType.FLOATDOUBLE→TSDataType.DOUBLEBOOLEAN→TSDataType.BOOLEANSTRING→TSDataType.TEXT其余未列出的类型会抛出UNSUPPORTED_DATA_TYPE异常。在表模型table模式下RelationalSeaTunnelRowSerializer.java 还额外支持TIMESTAMP与DATE字段值向 IoTDB 写入前的转换分别转为 epoch 毫秒与字符串。Sink 选项全解以下为 IoTDBv2 Sink 的全部选项继承自原文档并结合源码补充了默认值依据名称类型是否必填默认值描述node_urlsArray是-IoTDB 集群地址格式为[host1:port]或[host1:port,host2:port]usernameString是-IoTDB 用户名passwordString是-IoTDB 用户密码sql_dialectString否treeIoTDB 模型可选值为tree和table。tree表示树模型table表示表模型storage_groupString是-IoTDB 树模型指定设备路径前缀。例如设备路径为storage_group . key_device。如果key_device已经是完整设备路径可以设置为空字符串IoTDB 表模型指定数据库key_deviceString是-IoTDB 树模型在 SeaTunnelRow 中指定 IoTDB 设备 ID 的字段名IoTDB 表模型在 SeaTunnelRow 中指定 IoTDB 表名的字段名key_timestampString否数据处理时间IoTDB 树模型在 SeaTunnelRow 中指定 IoTDB 时间戳的字段名如未指定则使用处理时间作为时间戳IoTDB 表模型在 SeaTunnelRow 中指定 IoTDB 时间列的字段名如未指定则使用处理时间作为时间戳key_measurement_fieldsArray否见描述IoTDB 树模型在 SeaTunnelRow 中指定 IoTDB 测量列表的字段名如未指定则包括排除key_device与key_timestamp后的其余字段IoTDB 表模型在 SeaTunnelRow 中指定 IoTDB 测点列FIELD的字段名如未指定则包括排除key_device、key_timestamp、key_tag_fields、key_attribute_fields后的其余字段key_tag_fieldsArray否-IoTDB 树模型不生效IoTDB 表模型在 SeaTunnelRow 中指定 IoTDB 标签列TAG的字段名key_attribute_fieldsArray否-IoTDB 树模型不生效IoTDB 表模型在 SeaTunnelRow 中指定 IoTDB 属性列ATTRIBUTE的字段名batch_sizeInteger否1024缓存的记录数达到batch_size时连接器会把数据刷新到 IoTDB在 checkpoint 提交前和写入器关闭时也会刷新max_retriesInteger否0刷新失败时的最大重试次数retry_backoff_multiplier_msInteger否0计算重试等待时间的退避倍数单位为毫秒max_retry_backoff_msInteger否0最大重试等待时间单位为毫秒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 常用选项两种模型的字段语义速记树模型key_device用作 IoTDB 设备路径key_measurement_fields决定哪些字段写成测点。未配置key_measurement_fields时除key_device和key_timestamp之外的所有字段都会写成测点。表模型storage_group表示数据库key_device表示目标表名字段key_tag_fields表示 TAG 列key_attribute_fields表示 ATTRIBUTE 列key_measurement_fields表示 FIELD 列。未配置key_measurement_fields时除表名、时间、TAG 和 ATTRIBUTE 字段之外的所有字段都会写成 FIELD 列。选项的源码定义与加载所有选项均定义在 IoTDBv2SinkOptions.java继承自 IoTDBv2CommonOptions.javanode_urls、username、password、sql_dialect属于公共选项其中sql_dialect的默认值在源码中明确为tree即不配置时按树模型处理batch_size的默认值为1024key_device、storage_group在源码中为noDefaultValue()的必填项key_timestamp、key_measurement_fields、key_tag_fields、key_attribute_fields等均无默认值未配置即使用排除式的默认行为配置加载集中在 SinkConfig.java 的loadConfig(ReadonlyConfig)中完成zone_id会被解析为java.time.ZoneId对象。另外IoTDBv2Sink.java 的createWriter方法会依据sql_dialect是否为table常量定义于 SinkConstants.java分流表模型使用IoTDBv2RelationalSinkWriter其余情况含默认的tree使用IoTDBv2SinkWriter。示例 1写入 IoTDB 树模型数据树模型是 IoTDB 的传统数据组织方式root.存储组.设备.测点。以下示例基于 FakeSource 构造上游数据env { 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.0案例 1最小配置只填写所需的配置效果如下使用当前处理时间作为时间戳测点包括排除了key_device后的其余字段。sink { IoTDBv2 { node_urls [localhost:6667] username root password root key_device device_name # specify the deviceId use device_name field } }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| -----------------------------------------------------------------------------------------------------------------------------------------------------------------注意未配置key_timestamp时写入时间戳为 SeaTunnel 的处理时间源码中createTimestampExtractor在timestampKey为空时返回System.currentTimeMillis()。案例 2使用源事件的时间配置效果如下使用指定字段event_ts作为时间戳测点包括排除了key_device和key_timestamp后的其余字段。sink { IoTDBv2 { node_urls [localhost:6667] username root password root key_device device_name # specify the deviceId use device_name field key_timestamp event_ts # specify the timestamp use event_ts field } }IoTDB 中查询到的数据如下时间戳来自源数据event_ts 1664035200001IoTDB 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| -----------------------------------------------------------------------------------------------------------------------------------------------------------------案例 3使用源事件的时间和限定测量字段配置效果如下使用指定字段event_ts作为时间戳测点仅包括key_measurement_fields指定的temperature、moisture字段。sink { IoTDBv2 { node_urls [localhost:6667] username root password root key_device device_name key_timestamp event_ts key_measurement_fields [temperature, moisture] } }IoTDB 中查询到的数据如下只有两个测点列被写入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| -------------------------------------------------------------------------树模型下的设备路径拼接细节源码 DefaultSeaTunnelRowSerializer.java 的createDeviceExtractor会按以下规则生成最终设备路径若storage_group未配置空直接使用key_device字段值作为完整设备路径本示例即属此情况device_name本身已是root.test_group.device_a这样的完整路径若storage_group配置了值则拼接为storage_group . key_device字段值若storage_group以.结尾或设备值以.开头则直接拼接避免出现重复的.。同理createMeasurements在未配置key_measurement_fields时会默认取除deviceKey和timestampKey之外的所有字段作为测点列表。示例 2写入 IoTDB 表模型数据表模型Table Model是 IoTDB 2.x 引入的面向数据库/表/列的建模方式。将sql_dialect设置为table即可启用。以下示例基于 FakeSource 构造上游数据env { parallelism 2 job.mode BATCH } source { FakeSource { ... schema { fields { ts timestamp model_id string region string tag string status boolean arrival_date date temperature double } } } }上游 SeaTunnelRow 数据格式如下tsmodel_idregiontagstatusarrival_datetemperature2025-07-30T17:52:34.851id10700HKtag1true2024-11-124.342025-07-29T17:51:34.851id20700HKtag2false2024-12-015.542025-07-28T17:50:34.851id30700HKtag3false2024-12-227.34案例 1最小配置配置效果如下使用当前处理时间作为时间列测量列FIELD包括排除了key_device后的其余字段。sink { IoTDBv2 { node_urls [localhost:6667] username root password root sql_dialect table storage_group test_database key_device region } }此配置下region字段值如0700HK被用作表名写入到test_database数据库下。IoTDB 中的查询结果如下IoTDB SELECT * FROM test_database.0700HK; --------------------------------------------------------------------------------------------- | time| ts|model_id| tag|status|arrival_date|temperature| --------------------------------------------------------------------------------------------- |2025-08-14T17:52:34.85108:00|2025-07-30T17:52:34.851| id1|tag1| true| 2024-11-12| 4.34| |2025-08-14T17:51:34.85108:00|2025-07-29T17:51:34.851| id2|tag2| false| 2024-12-01| 5.54| |2025-08-14T17:50:34.85108:00|2025-07-28T17:50:34.851| id3|tag3| false| 2024-12-22| 7.34| ---------------------------------------------------------------------------------------------IoTDB DESC test_database.0700HK; ----------------------------- | ColumnName| DataType|Category| ----------------------------- | time|TIMESTAMP| TIME| | ts|TIMESTAMP| FIELD| | model_id| STRING| FIELD| | tag| STRING| FIELD| | status| BOOLEAN| FIELD| |arrival_date| DATE| FIELD| | temperature| DOUBLE| FIELD| -----------------------------案例 2使用源事件的时间和指定标签列、属性列配置效果如下使用指定字段ts作为时间列使用指定字段作为标签列TAGtag及属性列ATTRIBUTEmodel_id测量列FIELD包括排除了key_device、key_timestamp、key_tag_fields和key_attribute_fields后的其余字段即status、arrival_date、temperature。sink { IoTDBv2 { node_urls [localhost:6667] username root password root sql_dialect table storage_group test_database key_device region key_timestamp ts key_tag_fields [tag] key_attribute_fields [model_id] } }IoTDB 中的查询结果如下IoTDB SELECT * FROM test_database.0700HK; ---------------------------------------------------------------------- | time| tag|model_id|status|arrival_date|temperature| ---------------------------------------------------------------------- |2025-07-30T17:52:34.85108:00|tag1| id1| true| 2024-11-12| 4.34| |2025-07-29T17:51:34.85108:00|tag2| id2| false| 2024-12-01| 5.54| |2025-07-28T17:50:34.85108:00|tag3| id3| false| 2024-12-22| 7.34| ----------------------------------------------------------------------IoTDB DESC test_database.0700HK; ------------------------------ | ColumnName| DataType| Category| ------------------------------ | time|TIMESTAMP| TIME| | tag| STRING| TAG| | model_id| STRING|ATTRIBUTE| | status| BOOLEAN| FIELD| |arrival_date| DATE| FIELD| | temperature| DOUBLE| FIELD| ------------------------------可见tag列被标记为TAG类别model_id列被标记为ATTRIBUTE类别其余列均为FIELD。案例 3使用源事件的时间和限定测量列配置效果如下使用指定字段ts作为时间列使用key_measurement_fields指定测点列FIELD为status、temperature。sink { IoTDBv2 { node_urls [localhost:6667] username root password root sql_dialect table storage_group test_database key_device region key_timestamp ts key_measurement_fields [status, temperature] } }IoTDB 中的查询结果如下IoTDB SELECT * FROM test_database.0700HK; ---------------------------------------------- | time|status|temperature| ---------------------------------------------- |2025-07-30T17:52:34.85108:00| true| 4.34| |2025-07-29T17:51:34.85108:00| false| 5.54| |2025-07-28T17:50:34.85108:00| false| 7.34| ----------------------------------------------IoTDB DESC test_database.0700HK; ---------------------------- | ColumnName| DataType|Category| ---------------------------- | time|TIMESTAMP| TIME| | status| BOOLEAN| FIELD| |temperature| DOUBLE| FIELD| ---------------------------底层写入链路与实现原理Sink 创建与双模型分流在 IoTDBv2Sink.java 中createWriter依据sql_dialect是否为table创建不同的写入器table→IoTDBv2RelationalSinkWriter配合 RelationalSeaTunnelRowSerializer.java其他情况含默认值tree→IoTDBv2SinkWriter配合 DefaultSeaTunnelRowSerializer.java。写入流程以树模型为例见 IoTDBv2SinkWriter.java为write(SeaTunnelRow)→serializer.serialize(row)得到IoTDBv2Record→sinkClient.write(record)在prepareCommit()阶段会先执行sinkClient.flush()将缓存数据刷入存储后再快照状态。时间戳的提取规则无论树模型还是表模型时间戳提取逻辑一致源码createTimestampExtractor未配置key_timestamp或字段值为null时使用System.currentTimeMillis()即处理时间字段类型为STRING时按Long.parseLong解析为 epoch 毫秒字段类型为TIMESTAMP时按 UTC 时区转换为 epoch 毫秒字段类型为BIGINT时直接作为 epoch 毫秒使用其他类型会抛出UNSUPPORTED_DATA_TYPE。客户端批量写入与重试写入客户端 IoTDBv2SinkClient.java 承担了连接建立、批量缓冲、重试与刷新会话构建使用Session.Builder传入node_urls、username、password并按配置叠加thriftDefaultBufferSize、thriftMaxFrameSize、zoneIdsession.open()支持按enable_rpc_compression与connection_timeout_in_ms组合调用。批量缓冲每条记录先进入batchList当batchList.size() batch_size时触发flush()。批量写入flush()将缓存记录组装为BatchRecords调用 IoTDB 会话的insertRecords带类型列表或insertRecords纯字符串列表一次性写入。重试与退避写入失败时最多重试max_retries次每次等待时间为min(retry_backoff_multiplier_ms * i, max_retry_backoff_ms)毫秒超过重试上限则抛出FLUSH_DATA_FAILED。刷新时机除达到batch_size外prepareCommit()checkpoint 提交前与写入器close()时都会强制flush()保证数据完整落库。相关源码索引选项定义IoTDBv2SinkOptions.java、IoTDBv2CommonOptions.java配置加载SinkConfig.javaSink 入口与 Writer 分流IoTDBv2Sink.java、IoTDBv2SinkWriter.java、IoTDBv2RelationalSinkWriter.java序列化器DefaultSeaTunnelRowSerializer.java、RelationalSeaTunnelRowSerializer.java写入客户端IoTDBv2SinkClient.java、IoTDBv2RelationalSinkClient.java变更日志IoTDBv2 连接器的历史变更记录请参考 IoTDB 连接器变更日志。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考