☰
SeaTunnel Druid Sink Connector 使用指南:配置、数据类型映射与批量写入原理
2026/9/29 20:28:04 网站建设 项目流程
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel 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-once:Druid Sink 未实现两阶段提交或目标端去重机制,写入采用"批量提交任务"模型,因此在故障恢复场景下无法保证每条数据恰好只写入一次,适用于对一致性要求为 at-least-once 的场景。
  • support multiple table write:DruidSink实现了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 实现
TINYINTLONGLongDimensionSchema
SMALLINTLONGLongDimensionSchema
INTLONGLongDimensionSchema
BIGINTLONGLongDimensionSchema
FLOATFLOATFloatDimensionSchema
DOUBLEDOUBLEDoubleDimensionSchema
DECIMALDOUBLEDoubleDimensionSchema
STRINGSTRINGStringDimensionSchema
BOOLEANSTRINGStringDimensionSchema
TIMESTAMPSTRINGStringDimensionSchema

几点需要注意:

  • 整数类型(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否10000
common-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",但选项表格与源码中定义的默认值均为10000(DruidConfig.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,则可省略这两个参数。

配置示例

简单示例

来自官方文档的最小配置,仅写入单个 datasource:

sink { 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/sql,SELECT * 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_2),Sink 侧用${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);
  • 构造ParallelIndexSupervisorTask(Druid 原生批处理并行索引任务,可同时运行多个索引子任务,参见源码注释 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. 关闭时的最终 flush(close 方法)

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.1,HTTP 客户端为 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 引擎。

Changelog

next version

  • Add Druid sink connector(新增 Druid Sink 连接器)

参考资料(仓库内可深入阅读):

  • 官方文档:docs/en/connector-v2/sink/Druid.md
  • 配置项定义:DruidConfig.java
  • 连接器主类:DruidSink.java
  • 写入实现:DruidWriter.java
  • 工厂类:DruidSinkFactory.java
  • E2E 测试:DruidIT.java 及其 单表配置 / 多表配置
  • 相关概念:Sink Common Options、Connector V2 特性说明、Sink Options Placeholders
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询