- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
本篇技术指南以官方文档 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 Type | Druid Data Type | 源码中的维度 Schema 实现 |
|---|---|---|
| TINYINT | LONG | LongDimensionSchema |
| SMALLINT | LONG | LongDimensionSchema |
| INT | LONG | LongDimensionSchema |
| BIGINT | LONG | LongDimensionSchema |
| FLOAT | FLOAT | FloatDimensionSchema |
| DOUBLE | DOUBLE | DoubleDimensionSchema |
| DECIMAL | DOUBLE | DoubleDimensionSchema |
| STRING | STRING | StringDimensionSchema |
| BOOLEAN | STRING | StringDimensionSchema |
| TIMESTAMP | STRING | StringDimensionSchema |
几点需要注意:
- 整数类型(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:
| 名称 | 类型 | 是否必填 | 默认值 |
|---|---|---|---|
| coordinatorUrl | string | 是 | 无 |
| datasource | string | 是 | 无 |
| batchSize | int | 否 | 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 is
1024",但选项表格与源码中定义的默认值均为10000(DruidConfig.BATCH_SIZE_DEFAULT = 10000,见 DruidConfig.java)。请以 10000 为准,1024是文档笔误。
batchSize的取值需要权衡:值越小提交越频繁,数据可见性越好,但会产生更多索引任务与协调开销;值越大则 Druid 端批量索引效率更高,但内存缓冲与延迟相应增大。
common options(可选)
Sink 插件通用参数,主要用于多表写入场景下的数据流路由,详见 Sink Common Options:
| 名称 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| source_table_name | String | 否 | - | 指定当前插件处理哪个上游数据集(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.
相关推荐
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),仅供参考