SeaTunnel DuckDB Sink:通过 JDBC 连接器将数据写入 DuckDB 的完整配置指南
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文基于 SeaTunnel 官方文档 docs/en/connectors/sink/DuckDB.md,讲解如何通过 JDBC 连接器(DuckDB 方言)向 DuckDB 数据库文件写入数据,覆盖支持版本、依赖放置、完整 Sink 参数、数据类型映射与三类典型作业配置(简单写入、CDC 事件写入、Exactly-Once 写入)。读完本文后,你可以复制文中的 HOCON 示例直接配置任务,并理解每个参数在源码中的实际作用(方言识别、类型转换、Catalog 建表等),以便排查写入失败、类型不匹配等常见问题。
概述
SeaTunnel 的 DuckDB 写入能力由JDBC 连接器 + DuckDB 方言(dialect)实现,并非独立的 DuckDB 连接器。通过 JDBC 将数据写入 DuckDB 数据库文件,支持批处理与流处理两种模式,支持并发写入;当底层 JDBC 驱动提供 XA 数据源时(设置is_exactly_once = true并提供xa_data_source_class_name),还支持 Exactly-Once 语义。由于 DuckDB 是进程内(in-process)运行的嵌入式数据库,连接器面向的是本地数据库文件路径(jdbc:duckdb:/path/to/database.db)或内存数据库(jdbc:duckdb:)。
从源码结构看,方言识别的入口在 DuckDBDialectFactory:它通过@AutoService(JdbcDialectFactory.class)注册,并依据url.startsWith("jdbc:duckdb:")判断是否接管该连接,命中后创建 DuckDBDialect。因此只要配置url = "jdbc:duckdb:...",JDBC 连接器就会自动切换到 DuckDB 的方言行为(标识符引用、行转换、类型映射、表路径解析等)。
支持的 DuckDB 版本
- 0.8.x / 0.9.x / 0.10.x / 1.x
支持的引擎
Spark Flink SeaTunnel Zeta
依赖配置
依赖要求按引擎区分:
Spark / Flink 引擎
- 确保
duckdb_jdbc驱动 jar 包(Maven 坐标org.duckdb:duckdb_jdbc)已放置在${SEATUNNEL_HOME}/plugins/目录下。
SeaTunnel Zeta 引擎
- 确保
duckdb_jdbc驱动 jar 包已放置在${SEATUNNEL_HOME}/lib/目录下。
注意:DuckDB 是嵌入式数据库,JDBC 驱动包含本地库文件,驱动 jar 必须能被运行任务的进程加载,否则会出现ClassNotFoundException或驱动加载失败。
核心特性(Key Features)
| 特性 | 是否支持 | 说明 |
|---|---|---|
| exactly-once | 支持 | 使用 XA 事务保证 Exactly-Once。只有支持 XA 事务的数据库才能使用该语义,需设置is_exactly_once=true |
| cdc | 支持 | 可接收上游 CDC 事件流并写入 DuckDB |
| timer flush | 不支持 | 未勾选该特性 |
数据源信息
| 数据源 | 支持版本 | 驱动类 | URL 格式 | 获取方式 |
|---|---|---|---|---|
| DuckDB | 不同依赖版本对应不同驱动类 | org.duckdb.DuckDBDriver | jdbc:duckdb:/path/to/database.db | Maven 坐标org.duckdb:duckdb_jdbc |
数据类型映射
| SeaTunnel 数据类型 | DuckDB 数据类型 |
|---|---|
| BOOLEAN | BOOLEAN |
| TINYINT / SMALLINT / INT | INTEGER |
| BIGINT | BIGINT |
| DECIMAL(x,y)(列指定精度 < 38) | DECIMAL(x,y) |
| DECIMAL(x,y)(列指定精度 > 38) | DECIMAL(38,18) |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| STRING | VARCHAR |
| DATE | DATE |
| TIME | TIME |
| TIMESTAMP | TIMESTAMP |
| BYTES / ARRAY / ROW / MAP | BLOB |
类型转换的源码实现
上表对应 Sink 侧写入时的类型定义,核心实现在 DuckDBTypeConverter(通过@AutoService(TypeConverter.class)注册,identifier()返回DUCKDB)。其中两个值得注意的实现细节:
- DECIMAL 精度钳制:源码中定义了
MAX_PRECISION = 38、MAX_SCALE = 38,当精度超过 38 时会被截断到 38 并打印告警日志。这与上表“精度 > 38 映射为 DECIMAL(38,18)”的文档约定一致,配置高精度 DECIMAL 时应预期这种截断行为。 - 复杂类型的兜底:读取方向(Catalog 内省)中,DuckDB 的
ARRAY、STRUCT、MAP等复杂类型会降级为STRING并输出Complex type {} mapped to STRING, consider using JSON serialization的告警;未知类型同样回退为STRING。因此从 DuckDB 读数据时,复杂列建议改用 JSON 序列化后再同步。
写入时的行数据转换由 DuckDBJdbcRowConverter 完成,它继承AbstractJdbcRowConverter并按 DuckDB 方言执行 JDBCsetXxx绑定。
Sink 参数
| 参数 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL。示例:jdbc:duckdb:/path/to/database.db。内存数据库使用jdbc:duckdb: |
| driver | String | 是 | - | 连接远端数据源使用的 JDBC 驱动类名。DuckDB 固定为org.duckdb.DuckDBDriver |
| username | String | 否 | - | 连接用户名。DuckDB 本地文件无需认证,留空即可(除非你用自定义认证器包装) |
| password | String | 否 | - | 连接密码。DuckDB 本地文件无需认证,留空即可 |
| query | String | 否 | - | 直接指定写入 SQL(如INSERT ...)。设置query时优先于database/table/table_list |
| database | String | 否 | main | 与table配合自动生成 SQL 写入。与query互斥且优先级更高 |
| table | String | 否 | - | 与database配合自动生成 SQL 写入。与query互斥且优先级更高 |
| primary_keys | Array | 否 | - | 自动生成 SQL 时用于支持insert、delete、update等操作的主键字段 |
| connection_check_timeout_sec | Int | 否 | 30 | 校验连接所用数据库操作的等待超时时间(秒) |
| max_retries | Int | 否 | 0 | 失败的executeBatch调用的重试次数 |
| batch_size | Int | 否 | 1000 | 批量写入阈值:缓冲区记录数达到batch_size或时间达到checkpoint.interval时,将数据刷新入库 |
| is_exactly_once | Boolean | 否 | false | 是否启用 Exactly-Once 语义(基于 XA 事务)。启用时必须同时设置xa_data_source_class_name |
| generate_sink_sql | Boolean | 否 | false | 基于目标表结构生成 SQL 语句。需要配置database和table(或table_list) |
| xa_data_source_class_name | String | 否 | - | 数据库驱动的 XA 数据源类名。DuckDB 使用org.duckdb.DuckDBXADataSource |
| max_commit_attempts | Int | 否 | 3 | 事务提交失败的重试次数 |
| transaction_timeout_sec | Int | 否 | -1 | 事务打开后的超时时间,默认-1(永不超时)。注意:设置超时可能影响 Exactly-Once 语义 |
| auto_commit | Boolean | 否 | true | 是否自动提交事务。is_exactly_once = true时应设为false |
| field_ide | String | 否 | - | 字段名转换策略:ORIGINAL不转换;UPPERCASE转大写;LOWERCASE转小写 |
| properties | Map | 否 | - | 附加连接参数。当 properties 与 URL 存在同名参数时,优先级由驱动实现决定;对 DuckDB,properties 优先于 URL |
| common-options | - | 否 | - | Sink 插件通用参数 |
| schema_save_mode | Enum | 否 | CREATE_SCHEMA_WHEN_NOT_EXIST | 任务开始前如何处理目标端已有表结构。可选:RECREATE_SCHEMA、CREATE_SCHEMA_WHEN_NOT_EXIST、ERROR_WHEN_SCHEMA_NOT_EXIST |
| data_save_mode | Enum | 否 | APPEND_DATA | 任务开始前如何处理目标端已有数据。可选:DROP_DATA、APPEND_DATA、CUSTOM_PROCESSING、ERROR_WHEN_DATA_EXISTS |
| custom_sql | String | 否 | - | 当data_save_mode = CUSTOM_PROCESSING时填写,作为同步任务开始前执行的 SQL |
| enable_upsert | Boolean | 否 | true | 按primary_keys启用 upsert。若任务只有insert,设为false可加速数据导入 |
| multi_table_sink_replica | Int | 否 | 1 | 多表写入副本数。multi_table_sink_replica > 1时并行写入多个表 |
参数默认值与 JdbcCommonOptions 中的定义一致,例如
connection_check_timeout_sec默认 30 秒。该文件中url还声明了回退键base-url,即部分场景下写base-url也能被解析。
表名与标识符处理的方言细节
从源码结构看,DuckDB 方言对表路径有一套特殊的解析规则(DuckDBDialect):
parse(tablePath):形如a.b的表名解析为「库=default,schema=a,表=b」;只有一个片段的表名b解析为「库=default,schema=main,表=b」。这解释了为什么文档中database默认值为main(DuckDB 默认 schema)。quoteIdentifier使用双引号包裹标识符;tableIdentifier(TablePath)生成的写入目标形如"main"."sink_table"。hashModForField使用MOD(ABS(HASH(field)), mod)生成取模分片表达式,供partition_column并行分片使用。
关于 upsert 的一个重要说明
文档中enable_upsert默认值为true,但源码中DuckDBDialect.getUpsertStatement当前返回Optional.empty(),其注释说明:该连接器有意不提供行级 UPSERT SQL,因为 SeaTunnel 面向批处理 ETL 与追加写入场景,行级 UPSERT 在分析型存储引擎上可能造成显著性能退化。可以推断:对 DuckDB 而言,primary_keys更多用于生成 SQL 的辅助标识,CDC 事件的删除/更新落地行为依赖驱动层面的批量写入而非数据库侧 UPSERT 语句。如果你的任务只有纯插入且不需要主键处理,建议显式设置enable_upsert = false以提升导入速度(与文档 Tips 一致)。
Catalog:建表与表结构内省
除 Sink 写入外,仓库还提供了 DuckDB 的 Catalog 实现,支撑“目标表不存在时自动建表”(schema_save_mode = CREATE_SCHEMA_WHEN_NOT_EXIST)等能力:
- DuckDBCatalog 通过查询
information_schema.columns(并 LEFT JOINduckdb_columns获取列注释)内省列元数据;对无显式精度/标度的DECIMAL/NUMERIC列做确定性归一化(精度缺省补 38、负标度归 0),并处理表达式默认值(如CURRENT_TIMESTAMP)。其类注释明确指出:DuckDB 在 JVM 内对同一数据库存在“单连接”约束,该 Catalog 会集中管理并持有这条 JDBC 连接。 - 建表 SQL 由 DuckDBCreateTableSqlBuilder 基于
CatalogTable的列定义、主键与约束生成,字段名同样受field_ide策略影响。 - DuckDBCatalogFactory 与 DuckDBURLParser 负责按
jdbc:duckdb:URL 解析出文件路径信息。由于 DuckDB 无需用户名密码,JdbcCommonOptions 的注释也说明:不需要认证的数据源(如 DuckDB)应定义自己的optionRule(),而不套用baseCatalogRule()(后者强制要求 username/password)。
单元测试方面,可参考 DuckDBDialectTest、DuckDBTypeConverterTest 与 DuckDBCatalogTest 验证方言、类型转换与 Catalog 行为;DuckDBConnectDryRunValidationTest 覆盖了连接 dry-run 校验场景。
并发提示(Tips)
若未设置
partition_column,任务将以单并发运行;若设置了partition_column,则会按任务并发度并行执行(底层即上文hashModForField的MOD(ABS(HASH(column)), mod)分片逻辑)。
作业示例
示例一:简单写入(Simple)
批量模式,从 FakeSource 产生 1000 行数据写入本地 DuckDB 文件:
env { parallelism = 1 job.mode = "BATCH" } source { FakeSource { parallelism = 1 row_num = 1000 schema = { fields { id = "int" name = "string" age = "int" email = "string" } } } } sink { Jdbc { url = "jdbc:duckdb:/tmp/test.db" driver = "org.duckdb.DuckDBDriver" table = "sink_table" username = "" password = "" } }示例二:CDC(Change Data Capture)事件写入
流式模式,从 MySQL-CDC 捕获变更并写入 DuckDB。要点是开启generate_sink_sql并同时配置database与table,用primary_keys声明主键:
env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { MySQL-CDC { base-url = "jdbc:mysql://localhost:3306/test" username = "root" password = "123456" table-names = ["test.user"] } } sink { Jdbc { url = "jdbc:duckdb:/tmp/test.db" driver = "org.duckdb.DuckDBDriver" table = "sink_table" username = "" password = "" generate_sink_sql = true # You need to configure both database and table database = main table = "sink_table" primary_keys = ["id"] } }示例三:Exactly-Once 写入
通过 XA 事务保证精确一次语义。DuckDB 驱动提供的 XA 数据源类为org.duckdb.DuckDBXADataSource,启用时还需配合checkpoint.interval(流式)使事务在 checkpoint 处提交:
env { parallelism = 1 job.mode = "BATCH" } source { FakeSource { parallelism = 1 row_num = 1000 schema = { fields { id = "int" name = "string" age = "int" email = "string" } } } } sink { Jdbc { url = "jdbc:duckdb:/tmp/test.db" driver = "org.duckdb.DuckDBDriver" table = "sink_table" username = "" password = "" is_exactly_once = "true" xa_data_source_class_name = "org.duckdb.DuckDBXADataSource" } }使用注意事项小结
- URL 决定方言:
jdbc:duckdb:前缀会触发DuckDBDialectFactory.acceptsURL,无需额外指定方言;若使用其他 URL 格式,可结合dialect参数(见 JdbcCommonOptions)显式指定。 - 默认 schema 是 main:不写
database时默认使用main;单段表名会被解析到default库 +mainschema,跨 schema 写入请用两段式表名。 - DECIMAL 精度上限 38:超出的精度会被源码钳制为 38 并告警,配置上游高精度 DECIMAL 前请确认目标表能容纳。
- Embedded 单连接约束:DuckDBCatalog 的注释说明同一数据库在 JVM 内仅允许单条连接,多任务并发写同一 .db 文件时应让 DuckDB 侧开启多进程/多 worker(外部能力),而不要在单进程内假设多个独立连接。
- 复杂类型:BYTES/ARRAY/ROW/MAP 写入目标为 BLOB;从 DuckDB 读取复杂类型时会降级为 STRING(JSON 序列化推荐)。
Changelog
DuckDB 写入能力随 JDBC 连接器迭代,变更记录参见 connector-jdbc changelog。
相关文档
- Sink Common Options
- Connector V2 特性说明(exactly-once / CDC / timer flush)
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考