SeaTunnel DuckDB Sink:通过 JDBC 连接器将数据写入 DuckDB 的完整配置指南
2026/9/16 17:20:36 网站建设 项目流程

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 引擎

  1. 确保duckdb_jdbc驱动 jar 包(Maven 坐标org.duckdb:duckdb_jdbc)已放置在${SEATUNNEL_HOME}/plugins/目录下。

SeaTunnel Zeta 引擎

  1. 确保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.DuckDBDriverjdbc:duckdb:/path/to/database.dbMaven 坐标org.duckdb:duckdb_jdbc

数据类型映射

SeaTunnel 数据类型DuckDB 数据类型
BOOLEANBOOLEAN
TINYINT / SMALLINT / INTINTEGER
BIGINTBIGINT
DECIMAL(x,y)(列指定精度 < 38)DECIMAL(x,y)
DECIMAL(x,y)(列指定精度 > 38)DECIMAL(38,18)
FLOATFLOAT
DOUBLEDOUBLE
STRINGVARCHAR
DATEDATE
TIMETIME
TIMESTAMPTIMESTAMP
BYTES / ARRAY / ROW / MAPBLOB

类型转换的源码实现

上表对应 Sink 侧写入时的类型定义,核心实现在 DuckDBTypeConverter(通过@AutoService(TypeConverter.class)注册,identifier()返回DUCKDB)。其中两个值得注意的实现细节:

  1. DECIMAL 精度钳制:源码中定义了MAX_PRECISION = 38MAX_SCALE = 38,当精度超过 38 时会被截断到 38 并打印告警日志。这与上表“精度 > 38 映射为 DECIMAL(38,18)”的文档约定一致,配置高精度 DECIMAL 时应预期这种截断行为。
  2. 复杂类型的兜底:读取方向(Catalog 内省)中,DuckDB 的ARRAYSTRUCTMAP等复杂类型会降级为STRING并输出Complex type {} mapped to STRING, consider using JSON serialization的告警;未知类型同样回退为STRING。因此从 DuckDB 读数据时,复杂列建议改用 JSON 序列化后再同步。

写入时的行数据转换由 DuckDBJdbcRowConverter 完成,它继承AbstractJdbcRowConverter并按 DuckDB 方言执行 JDBCsetXxx绑定。

Sink 参数

参数类型是否必填默认值描述
urlString-JDBC 连接 URL。示例:jdbc:duckdb:/path/to/database.db。内存数据库使用jdbc:duckdb:
driverString-连接远端数据源使用的 JDBC 驱动类名。DuckDB 固定为org.duckdb.DuckDBDriver
usernameString-连接用户名。DuckDB 本地文件无需认证,留空即可(除非你用自定义认证器包装)
passwordString-连接密码。DuckDB 本地文件无需认证,留空即可
queryString-直接指定写入 SQL(如INSERT ...)。设置query时优先于database/table/table_list
databaseStringmaintable配合自动生成 SQL 写入。与query互斥且优先级更高
tableString-database配合自动生成 SQL 写入。与query互斥且优先级更高
primary_keysArray-自动生成 SQL 时用于支持insertdeleteupdate等操作的主键字段
connection_check_timeout_secInt30校验连接所用数据库操作的等待超时时间(秒)
max_retriesInt0失败的executeBatch调用的重试次数
batch_sizeInt1000批量写入阈值:缓冲区记录数达到batch_size或时间达到checkpoint.interval时,将数据刷新入库
is_exactly_onceBooleanfalse是否启用 Exactly-Once 语义(基于 XA 事务)。启用时必须同时设置xa_data_source_class_name
generate_sink_sqlBooleanfalse基于目标表结构生成 SQL 语句。需要配置databasetable(或table_list
xa_data_source_class_nameString-数据库驱动的 XA 数据源类名。DuckDB 使用org.duckdb.DuckDBXADataSource
max_commit_attemptsInt3事务提交失败的重试次数
transaction_timeout_secInt-1事务打开后的超时时间,默认-1(永不超时)。注意:设置超时可能影响 Exactly-Once 语义
auto_commitBooleantrue是否自动提交事务。is_exactly_once = true时应设为false
field_ideString-字段名转换策略:ORIGINAL不转换;UPPERCASE转大写;LOWERCASE转小写
propertiesMap-附加连接参数。当 properties 与 URL 存在同名参数时,优先级由驱动实现决定;对 DuckDB,properties 优先于 URL
common-options--Sink 插件通用参数
schema_save_modeEnumCREATE_SCHEMA_WHEN_NOT_EXIST任务开始前如何处理目标端已有表结构。可选:RECREATE_SCHEMACREATE_SCHEMA_WHEN_NOT_EXISTERROR_WHEN_SCHEMA_NOT_EXIST
data_save_modeEnumAPPEND_DATA任务开始前如何处理目标端已有数据。可选:DROP_DATAAPPEND_DATACUSTOM_PROCESSINGERROR_WHEN_DATA_EXISTS
custom_sqlString-data_save_mode = CUSTOM_PROCESSING时填写,作为同步任务开始前执行的 SQL
enable_upsertBooleantrueprimary_keys启用 upsert。若任务只有insert,设为false可加速数据导入
multi_table_sink_replicaInt1多表写入副本数。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,则会按任务并发度并行执行(底层即上文hashModForFieldMOD(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并同时配置databasetable,用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" } }

使用注意事项小结

  1. URL 决定方言jdbc:duckdb:前缀会触发DuckDBDialectFactory.acceptsURL,无需额外指定方言;若使用其他 URL 格式,可结合dialect参数(见 JdbcCommonOptions)显式指定。
  2. 默认 schema 是 main:不写database时默认使用main;单段表名会被解析到default库 +mainschema,跨 schema 写入请用两段式表名。
  3. DECIMAL 精度上限 38:超出的精度会被源码钳制为 38 并告警,配置上游高精度 DECIMAL 前请确认目标表能容纳。
  4. Embedded 单连接约束:DuckDBCatalog 的注释说明同一数据库在 JVM 内仅允许单条连接,多任务并发写同一 .db 文件时应让 DuckDB 侧开启多进程/多 worker(外部能力),而不要在单进程内假设多个独立连接。
  5. 复杂类型: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),仅供参考

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

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

立即咨询