尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

SeaTunnel Druid Sink Connector 使用指南:配置、数据类型映射与批量写入原理

发布时间:2026/9/29 20:28:06

资讯中心
01
ARTICLE

SeaTunnel Druid Sink Connector 使用指南:配置、数据类型映射与批量写入原理

SeaTunnel Druid Sink Connector 使用指南:配置、数据类型映射与批量写入原理
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本篇技术指南以官方文档 Druid.md 为骨架结合connector-druid模块源码与 E2E 测试展开介绍如何将 SeaTunnel 中的数据写入 Apache Druid涵盖连接器能力、配置参数、数据类型映射、占位符多表写入以及底层基于 HTTP 任务提交的批量写入机制。连接器概述Druid Sink Connector 是 SeaTunnel Connector V2 体系中的一个写入端插件用于将上游数据如来自 FakeSource、Kafka、JDBC 等 Source 的数据批量写入 Apache Druid。Druid 是一款面向实时分析场景的列式数据仓库常用于 OLAP 查询与聚合分析因此该连接器适合把 SeaTunnel 作业清洗、转换后的结果数据导入 Druid 供后续分析查询使用。在仓库中该连接器位于seatunnel-connectors-v2/connector-druid模块插件注册名为Druid见 DruidSink.java 中getPluginName()的返回值。同时它也出现在发布相关配置中plugin-mapping.properties中映射seatunnel.sink.Druid connector-druid默认打包列表 plugin_config 也包含connector-druid说明在默认发行包中即可直接使用。核心特性支持情况根据官方特性文档 connector-v2-features.md 的定义Druid Sink 支持能力如下特性支持情况exactly-once精确一次❌ 不支持support multiple table write多表写入✅ 支持exactly-onceDruid Sink 未实现两阶段提交或目标端去重机制写入采用批量提交任务模型因此在故障恢复场景下无法保证每条数据恰好只写入一次适用于对一致性要求为 at-least-once 的场景。support multiple table writeDruidSink实现了SupportMultiTableSink接口见 DruidSink.java配合 sink-options-placeholders.md 中的占位符能力可以在单个作业中动态写入多个 Druid datasource。数据类型映射Druid 的维度Dimension数据类型与 SeaTunnel 类型并非一一对应。官方文档给出的映射表如下其底层逻辑在 DruidWriter.java 的transformToDimensionSchema()方法中得到了源码级印证——该方法根据 SeaTunnel 字段的SqlType将每个字段声明为 Druid 的DimensionSchema子类SeaTunnel Data TypeDruid Data Type源码中的维度 Schema 实现TINYINTLONGLongDimensionSchemaSMALLINTLONGLongDimensionSchemaINTLONGLongDimensionSchemaBIGINTLONGLongDimensionSchemaFLOATFLOATFloatDimensionSchemaDOUBLEDOUBLEDoubleDimensionSchemaDECIMALDOUBLEDoubleDimensionSchemaSTRINGSTRINGStringDimensionSchemaBOOLEANSTRINGStringDimensionSchemaTIMESTAMPSTRINGStringDimensionSchema几点需要注意整数类型TINYINT/SMALLINT/INT/BIGINT统一映射为 Druid 的LONG维度因此超出 64 位有符号整数范围的值无法安全写入DECIMAL会被降级为DOUBLE存在精度损失风险对精度敏感的场景需在写入前自行评估BOOLEAN与TIMESTAMP被映射为STRING维度其中 TIMESTAMP 字段在写入时以字符串形式落入 Druid若字段类型不在上述映射范围内如NULL、BYTES、ARRAY、MAP等transformToDimensionSchema()会抛出DruidConnectorException错误码UNSUPPORTED_DATA_TYPE作业将失败——这也是该连接器当前的类型边界。配置参数详解官方文档的参数表如下配置项定义可参见 DruidConfig.java名称类型是否必填默认值coordinatorUrlstring是无datasourcestring是无batchSizeint否10000common-options否-DruidSinkFactory的optionRule()将coordinatorUrl与datasource声明为必填项见 DruidSinkFactory.java缺少任一必填项都会在作业校验阶段报错。coordinatorUrl [string]必填Druid Coordinator协调服务的主机与端口格式为host:port。连接器会向该地址发送 HTTP 请求来提交索引任务。示例myHost:8888。在仓库的 E2E 测试 DruidIT.java 中测试环境通过 Docker Compose 启动 Druid对外暴露的正是router服务的8888端口因此文档示例中的端口8888与官方默认部署方式保持一致。datasource [string]必填要写入的 Druid datasource数据源名称可理解为 Druid 侧的一张分析表。示例seatunnel。batchSize [int]可选默认 10000每批次写入 Druid 的行数阈值。当累积行数达到batchSize时触发一次批量提交。⚠️ 文档正文中对该参数的描述文字写的是 Default value is1024但选项表格与源码中定义的默认值均为10000DruidConfig.BATCH_SIZE_DEFAULT 10000见 DruidConfig.java。请以 10000 为准1024是文档笔误。batchSize的取值需要权衡值越小提交越频繁数据可见性越好但会产生更多索引任务与协调开销值越大则 Druid 端批量索引效率更高但内存缓冲与延迟相应增大。common options可选Sink 插件通用参数主要用于多表写入场景下的数据流路由详见 Sink Common Options名称类型必填默认值说明source_table_nameString否-指定当前插件处理哪个上游数据集result_table_name对应的表不指定时处理配置文件中上一插件的输出数据注意事项摘自 common-options 文档当配置了source_table_name时上游必须同时设置result_table_name若作业中 source、transform、sink 任一环节的插件数量大于 1则必须为每个连接器显式指定source_table_name与result_table_name若作业只有一个 source、一个或零个transform、一个 sink则可省略这两个参数。配置示例简单示例来自官方文档的最小配置仅写入单个 datasourcesink { Druid { coordinatorUrl testHost:8888 datasource seatunnel } }占位符获取上游表元数据示例利用sink-options-placeholders能力datasource中可以引用上游 Catalog 表元数据实现一张表对应一个 datasource的动态映射sink { Druid { coordinatorUrl testHost:8888 datasource ${table_name}_test } }完整的批量作业示例来自 E2E 测试仓库中 fakesource_to_druid.conf 给出了一个可直接运行的完整作业配置覆盖了从 FakeSource 生成全类型数据到写入 Druid 的完整链路env { parallelism 1 job.mode BATCH } source { FakeSource { result_table_name fake schema { fields { c_boolean boolean c_timestamp timestamp c_string string c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double c_decimal decimal(16, 1) } } rows [ { kind INSERT fields [true, 2020-02-02T02:02:02, NEW, 1, 2, 3, 4, 4.3, 5.3, 6.3] }, { kind INSERT fields [false, 2012-12-21T12:34:56, AAA, 1, 1, 333, 323232, 3.1, 9.33333, 99999.99999999] } ] } } transform { } sink { Druid { coordinatorUrl localhost:8888 datasource testDataSource } }该配置覆盖了文档数据类型映射表中的全部 10 种 SeaTunnel 类型boolean、timestamp、string、tinyint、smallint、int、bigint、float、double、decimal可作为自测映射关系的参考。E2E 测试 DruidIT.java 通过 Druid SQL 查询接口POST /druid/v2/sqlSELECT * FROM testDataSource验证了写入结果与预期数据完全一致且布尔值以true/false字符串形式存储时间戳以2020-02-02T02:02:02字符串形式存储与映射表行为吻合。多表写入示例占位符 multi-table仓库中 fakesource_to_druid_with_multi.conf 展示了多表写入的标准写法FakeSource 通过tables_configs声明多张表druid_sink_1、druid_sink_2Sink 侧用${table_name}占位符动态解析表名env { parallelism 1 job.mode BATCH } source { FakeSource { tables_configs [ { schema { table druid_sink_1 fields { id int val_bool boolean val_tinyint tinyint val_smallint smallint val_int int val_bigint bigint val_float float val_double double val_decimal decimal(16, 1) val_string string } } rows [ { kind INSERT fields [1, true, 1, 2, 3, 4, 4.3, 5.3, 6.3, NEW] } ] }, { schema { table druid_sink_2 fields { id int val_bool boolean val_tinyint tinyint val_smallint smallint val_int int val_bigint bigint val_float float val_double double val_decimal decimal(16, 1) } } rows [ { kind INSERT fields [1, true, 1, 2, 3, 4, 4.3, 5.3, 6.3] } ] } ] } } transform { } sink { Druid { coordinatorUrl localhost:8888 datasource ${table_name} } }关于占位符的完整语法${database_name}、${schema_name}、${table_name}、${primary_key}、${field_names}等及默认值写法${table_name:default}请参考 sink-options-placeholders.md。需要注意占位符替换发生在连接器启动之前若上游表元数据缺少对应字段例如 MySQL 源没有schema_name、Oracle 源没有database_name占位符将不会被替换进而导致配置异常。另外多表读取能力在 Spark/Flink 引擎上暂不支持E2E 测试中对 SPARK、FLINK 引擎跳过了多表用例见 DruidIT.java请使用 SeaTunnel Zeta 引擎运行多表写入作业。底层写入原理从 SeaTunnel 行到 Druid 索引任务该连接器的核心实现在 DruidWriter.java写入流程可分为行缓冲 → 批量提交任务两个阶段。1. 行数据转 CSV 缓冲write 方法DruidWriter.write()DruidWriter.java将每条SeaTunnelRow的所有非空字段按,拼接为一行 CSV并自动在每行末尾追加一个timestamp字段其值取当前系统时间System.currentTimeMillis()即写入时刻的 processTime。这是 Druid 数据模型的硬性要求——Druid 要求每条记录必须包含主时间戳列primary timestamp源码注释也明确指出这一点见 DruidWriter.java。因此最终提交给 Druid 的 CSV 中列顺序为上游表的所有字段 自动追加的timestamp列。这意味着你在查询 Druid datasource 时会看到额外多出的一列timestamp其值来自数据被 SeaTunnel 写入的时刻而非业务时间。2. 批量触发与任务提交flush 方法当缓冲行数达到batchSize时currentBatchSize batchSize调用flush()DruidWriter.java执行批量提交使用InlineInputSource将内存中的 CSV 文本作为数据源配合CsvInputFormat声明列名上游字段 timestamp构造ParallelIndexSupervisorTaskDruid 原生批处理并行索引任务可同时运行多个索引子任务参见源码注释 DruidWriter.java其中DataSchema使用TimestampSpec(timestamp, auto, null)解析主时间戳并以UniformGranularitySpec(Granularities.HOUR, Granularities.MINUTE, false, null)设置分段粒度——segment 粒度为小时、查询粒度为分钟通过 Apache HttpClient 向http://{coordinatorUrl}/druid/indexer/v1/task发送POST请求Content-Type: application/json提交任务 JSON任务 JSON 在序列化后会被移除id、groupId、resource以及spec.tuningConfig等运行时字段见 DruidWriter.java仅保留可提交的最小任务定义。3. 关闭时的最终 flushclose 方法close()DruidWriter.java会先执行最后一次flush()确保缓冲中不足batchSize的残余数据也被提交然后关闭 HTTP 客户端。也就是说无论数据量大小作业结束时数据都会被提交到 Druid只是可能被拆分成多个批次任务。4. 调用链小结完整的调用链可概括为SeaTunnel 作业FakeSource/Kafka/JDBC 等 Source → DruidSink.createWriter() 创建 DruidWriter读取 coordinatorUrl/datasource/batchSize 配置 → DruidWriter.write(row) 逐行转 CSV 并缓存到达 batchSize 触发 flush → flush() 构造 ParallelIndexSupervisorTask → POST http://{coordinatorUrl}/druid/indexer/v1/task 提交批量索引任务 → Druid Overlord/Peons 完成数据落盘与 segment 构建 → close() 最终 flush 残余数据其中DruidSink.createWriter()见 DruidSink.java负责将配置项与上游SeaTunnelRowType来自 CatalogTable注入 Writer而该连接器依赖的 Druid 客户端库版本为 24.0.1HTTP 客户端为 httpclient 4.5.13见 connector-druid/pom.xml如果你自行部署的 Druid 版本与此差异较大需要关注任务 API 的兼容性。注意事项与使用建议版本兼容性connector-druid编译依赖 Druid 24.0.1 的相关库建议目标集群为相近版本E2E 测试中 Spark 2.4 容器因 RoaringBitmap 版本不兼容被禁用见 DruidIT.java在 Spark 2.4 上运行该连接器可能遇到类似依赖冲突。一致性语义该连接器不支持 exactly-once故障重启时可能出现重复写入请按 at-least-once 设计下游去重或幂等策略。timestamp 自动追加每行数据都会附带写入时刻的timestamp列这是 Druid 主时间戳要求所致查询结果中会多出该列。DECIMAL 精度decimal类型写入后为DOUBLE存在精度损失高精度场景需提前处理。批量提交即数据可见写入采用攒批提交索引任务模式数据在任务完成 segment 构建后才可被查询小批量会带来更多索引任务开销。多表写入使用${table_name}等占位符时请确认上游提供对应元数据且优先使用 SeaTunnel Zeta 引擎。Changelognext versionAdd Druid sink connector新增 Druid Sink 连接器参考资料仓库内可深入阅读官方文档docs/en/connector-v2/sink/Druid.md配置项定义DruidConfig.java连接器主类DruidSink.java写入实现DruidWriter.java工厂类DruidSinkFactory.javaE2E 测试DruidIT.java 及其 单表配置 / 多表配置相关概念Sink Common Options、Connector V2 特性说明、Sink Options Placeholders赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel JDBC Snowflake Sink Connector 使用指南配置、类型映射与 CDC 写入实践SeaTunnel JDBC Snowflake Sink Connector 使用指南配置、类型映射与 CDC 写入实践 SeaTunnel 的 Snowf数据工程大数据批处理流处理SeaTunnel Vertica Sink Connector 实战指南JDBC 写入、类型映射与 Exactly-Once 配置SeaTunnel Vertica Sink Connector 实战指南JDBC 写入、类型映射与 Exactly Once 配置 本文以仓库中 docs/数据工程大数据批处理流处理SeaTunnel IoTDB Sink 连接器实战指南配置、数据类型映射与写入原理SeaTunnel IoTDB Sink 连接器实战指南配置、数据类型映射与写入原理 SeaTunnel 的 IoTDB Sink 连接器用于将 SeaTun数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。