ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

SeaTunnel DuckDB Sink:通过 JDBC 连接器将数据写入 DuckDB 的完整配置指南

SeaTunnel DuckDB Sink:通过 JDBC 连接器将数据写入 DuckDB 的完整配置指南 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_oncetruecdc支持可接收上游 CDC 事件流并写入 DuckDBtimer flush不支持未勾选该特性数据源信息数据源支持版本驱动类URL 格式获取方式DuckDB不同依赖版本对应不同驱动类org.duckdb.DuckDBDriverjdbc:duckdb:/path/to/database.dbMaven 坐标org.duckdb:duckdb_jdbc数据类型映射SeaTunnel 数据类型DuckDB 数据类型BOOLEANBOOLEANTINYINT / SMALLINT / INTINTEGERBIGINTBIGINTDECIMAL(x,y)列指定精度 38DECIMAL(x,y)DECIMAL(x,y)列指定精度 38DECIMAL(38,18)FLOATFLOATDOUBLEDOUBLESTRINGVARCHARDATEDATETIMETIMETIMESTAMPTIMESTAMPBYTES / ARRAY / ROW / MAPBLOB类型转换的源码实现上表对应 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 参数参数类型是否必填默认值描述urlString是-JDBC 连接 URL。示例jdbc:duckdb:/path/to/database.db。内存数据库使用jdbc:duckdb:driverString是-连接远端数据源使用的 JDBC 驱动类名。DuckDB 固定为org.duckdb.DuckDBDriverusernameString否-连接用户名。DuckDB 本地文件无需认证留空即可除非你用自定义认证器包装passwordString否-连接密码。DuckDB 本地文件无需认证留空即可queryString否-直接指定写入 SQL如INSERT ...。设置query时优先于database/table/table_listdatabaseString否main与table配合自动生成 SQL 写入。与query互斥且优先级更高tableString否-与database配合自动生成 SQL 写入。与query互斥且优先级更高primary_keysArray否-自动生成 SQL 时用于支持insert、delete、update等操作的主键字段connection_check_timeout_secInt否30校验连接所用数据库操作的等待超时时间秒max_retriesInt否0失败的executeBatch调用的重试次数batch_sizeInt否1000批量写入阈值缓冲区记录数达到batch_size或时间达到checkpoint.interval时将数据刷新入库is_exactly_onceBoolean否false是否启用 Exactly-Once 语义基于 XA 事务。启用时必须同时设置xa_data_source_class_namegenerate_sink_sqlBoolean否false基于目标表结构生成 SQL 语句。需要配置database和table或table_listxa_data_source_class_nameString否-数据库驱动的 XA 数据源类名。DuckDB 使用org.duckdb.DuckDBXADataSourcemax_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务打开后的超时时间默认-1永不超时。注意设置超时可能影响 Exactly-Once 语义auto_commitBoolean否true是否自动提交事务。is_exactly_once true时应设为falsefield_ideString否-字段名转换策略ORIGINAL不转换UPPERCASE转大写LOWERCASE转小写propertiesMap否-附加连接参数。当 properties 与 URL 存在同名参数时优先级由驱动实现决定对 DuckDBproperties 优先于 URLcommon-options-否-Sink 插件通用参数schema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST任务开始前如何处理目标端已有表结构。可选RECREATE_SCHEMA、CREATE_SCHEMA_WHEN_NOT_EXIST、ERROR_WHEN_SCHEMA_NOT_EXISTdata_save_modeEnum否APPEND_DATA任务开始前如何处理目标端已有数据。可选DROP_DATA、APPEND_DATA、CUSTOM_PROCESSING、ERROR_WHEN_DATA_EXISTScustom_sqlString否-当data_save_mode CUSTOM_PROCESSING时填写作为同步任务开始前执行的 SQLenable_upsertBoolean否true按primary_keys启用 upsert。若任务只有insert设为false可加速数据导入multi_table_sink_replicaInt否1多表写入副本数。multi_table_sink_replica 1时并行写入多个表参数默认值与 JdbcCommonOptions 中的定义一致例如connection_check_timeout_sec默认 30 秒。该文件中url还声明了回退键base-url即部分场景下写base-url也能被解析。表名与标识符处理的方言细节从源码结构看DuckDB 方言对表路径有一套特殊的解析规则DuckDBDialectparse(tablePath)形如a.b的表名解析为「库defaultschemaa表b」只有一个片段的表名b解析为「库defaultschemamain表b」。这解释了为什么文档中database默认值为mainDuckDB 默认 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 } }示例二CDCChange 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 读取复杂类型时会降级为 STRINGJSON 序列化推荐。ChangelogDuckDB 写入能力随 JDBC 连接器迭代变更记录参见 connector-jdbc changelog。相关文档Sink Common OptionsConnector 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),仅供参考
返回列表