ARTICLE DETAIL

资讯详情

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

SeaTunnel JDBC DuckDB Source Connector 完全指南:从本地数据库文件到分布式数据管道的读取实战

SeaTunnel JDBC DuckDB Source Connector 完全指南:从本地数据库文件到分布式数据管道的读取实战 SeaTunnel JDBC DuckDB Source Connector 完全指南从本地数据库文件到分布式数据管道的读取实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 通过 JDBC Source Connector 支持读取 DuckDB 数据库中的数据。DuckDB 是一款进程内in-processSQL 分析型OLAP数据库没有独立的远程服务端连接器直接面向本地数据库文件或内存数据库工作。本指南将完整讲解该连接器的能力边界、依赖安装、全部 Source 配置参数、数据类型映射、并行切片读取原理并结合仓库源码给出可直接复制运行的实战配置。概述与适用场景DuckDB 以单文件数据库的形式存在通常用于本地分析、数据探索与小型数据管道。SeaTunnel 的 JDBC DuckDB Source 通过jdbc:duckdb:/path/to/database.db这样的 JDBC URL 直接打开本地数据库文件也可以连接内存数据库例如jdbc:duckdb:memory:。从连接器角度看DuckDB 场景有以下几个显著特点无远程服务器连接串指向本地文件路径或memory:不存在 host/port 概念批量读取友好官方特性表中batch、exactly-once、column projection、parallelism、user-defined split均为支持状态stream流式模式不适用嵌入式数据库本身没有持续变更日志可供消费查询即投影连接器支持自定义查询 SQL通过select指定列即可实现列裁剪projection效果多表一次作业读取通过table_list可以在一个 Source 中并行读取同一数据库文件内的多张表。该连接器的官方文档位于 docs/en/connectors/source/DuckDB.md底层实现归属于 JDBC 连接器模块相关源码集中在seatunnel-connectors-v2/connector-jdbc。支持版本与运行引擎维度支持情况DuckDB 版本0.8.x / 0.9.x / 0.10.x / 1.x运行引擎Spark、Flink、SeaTunnel Zeta提示不同依赖版本的 DuckDB 驱动类名可能不同配置前请先确认当前duckdb_jdbc驱动包版本对应的 Driver 类。官方文档给出的驱动类为org.duckdb.DuckDBDriver。使用依赖驱动包放置位置DuckDB 驱动 jarduckdb_jdbc可通过 Maven 中央仓库获取需要手动放置到 SeaTunnel 安装目录中放置位置取决于运行引擎Spark / Flink 引擎将驱动 jar 放入${SEATUNNEL_HOME}/plugins/目录SeaTunnel Zeta 引擎将驱动 jar 放入${SEATUNNEL_HOME}/lib/目录。放置完成后即可在作业配置中使用Jdbc插件声明 DuckDB 数据源。连接器能力清单以下是官方特性清单详见 connector-v2-features能力支持状态batch批量✅ 支持stream流式❌ 不支持exactly-once精确一次✅ 支持column projection列投影✅ 支持parallelism并行✅ 支持user-defined split用户自定义切片✅ 支持其中列投影能力的实现方式为支持自定义查询 SQL通过查询语句控制读取的字段集合。数据源信息速查数据源支持版本驱动URL 格式驱动获取DuckDB不同依赖版本驱动类可能不同org.duckdb.DuckDBDriverjdbc:duckdb:/path/to/database.dbMaven 中央仓库duckdb_jdbc构件从源码看JDBC 连接器通过 DuckDBDialectFactory 中的acceptsURL方法以url.startsWith(jdbc:duckdb:)识别 DuckDB 方言因此 URL 前缀jdbc:duckdb:是判别该数据源类型的硬性条件。Source 配置参数详解连接器在作业配置中的插件名为JdbcJDBC 连接器统一插件名。完整参数如下参数名类型是否必填默认值说明urlString是-JDBC 连接 URL例如jdbc:duckdb:/path/to/database.dbdriverString是-JDBC 驱动类名DuckDB 场景填org.duckdb.DuckDBDriverusernameString否-连接用户名passwordString否-连接密码queryString是-查询语句connection_check_timeout_secInt否30等待连接校验操作完成的超时时间秒partition_columnString否-并行切分的列名仅支持数值型主键且只能配置一个列partition_lower_boundBigDecimal否-扫描的partition_column最小值不设置时 SeaTunnel 会查询数据库自动获取 min 值partition_upper_boundBigDecimal否-扫描的partition_column最大值不设置时 SeaTunnel 会查询数据库自动获取 max 值partition_numInt否job parallelism切分数量仅支持正整数默认等于作业并行度fetch_sizeInt否0查询返回大量对象时可通过设置行抓取大小row fetch size减少数据库访问次数以提升性能0 表示使用 JDBC 默认值propertiesMap否-额外连接配置参数当 properties 与 URL 中参数同名时优先级由驱动具体实现决定DuckDB 中 properties 优先于 URLtable_pathString否-表的完整路径可用其替代query例如main.table1table_listArray否-待读取表列表可用其替代table_path例如[{ table_path main.table1 }, { table_path main.table2, query select id, name from main.table2 }]where_conditionString否-应用于所有表/查询的公共行过滤条件必须以where开头例如where id 100split.sizeInt否8096单次切片包含的行数读取表时会将表按此行数拆分为多个 splitcommon-options-否-Source 插件公共参数详见 Source Common Options参数实现细节与取值建议从源码 JdbcSourceOptions 中可以确认以下实现细节split.size默认值 8096定义于JdbcSourceOptions.java#L49-L54即每 8096 行切一个 splitfetch_size默认 0表示交由 JDBC 驱动使用默认抓取大小table_path、where_condition、table_list均为无默认值的可选参数其中table_list被定义为ListJdbcSourceTableConfig结构类型除文档列出的参数外JDBC 源还提供一组split.*高级调优参数例如split.even-distribution.factor.upper-bound默认 100.0与split.even-distribution.factor.lower-bound默认 0.05用于判定表数据分布是否均匀分布因子计算公式为(MAX(id) - MIN(id) 1) / rowCount均匀分布时走均匀切分优化不均匀时走查询式切分split.sample-sharding.threshold默认 1000、split.inverse-sampling.rate默认 1000、split.allow-sampling默认 true控制大数据量下基于采样的分片策略use_select_count默认 false、skip_analyze默认 false控制表行数的统计方式。这些参数同样适用于 DuckDB 数据源可在需要精细控制切片行为时使用。关于 table_list 的结构约束table_list中每一项都对应一张表从 JdbcSourceTableConfig 源码可见每项可独立配置table_path表完整路径必填项query该表自定义查询用于过滤行与列partition_column/partition_num/partition_lower_bound/partition_upper_bound该表的独立并行切分配置use_select_count/skip_analyze/use_regex该表的统计与正则匹配开关。需要注意两点实现约束当table_list中配置了多张表时各表的table_path必须唯一不允许为空或重复否则校验会直接抛异常见JdbcSourceTableConfig.java#L108-L119当表项未显式指定partition_num时会使用默认值10见JdbcSourceTableConfig.java#L42。数据类型映射DuckDB 类型到 SeaTunnel 类型的官方映射关系如下DuckDB 数据类型SeaTunnel 数据类型BOOLEANBOOLEANTINYINTTINYINTUTINYINTSMALLINTSMALLINTUSMALLINTINTEGERINTUINTEGERBIGINTBIGINTUBIGINTDECIMAL(20,0)HUGEINTDECIMAL(38,0)FLOATFLOATDOUBLEDOUBLEDECIMAL(x,y)列大小 38DECIMAL(x,y)DECIMAL(x,y)列大小 38DECIMAL(38,18)VARCHARCHARTEXTJSONUUIDINTERVALSTRINGDATEDATETIMETIMETIMESTAMPTIMESTAMP WITH TIME ZONETIMESTAMPBLOBARRAYSTRUCTMAPBYTES源码层面的类型转换细节类型映射的落地实现位于 DuckDBTypeConverter并通过 DuckDBTypeMapper 接入 JDBC 方言体系。从当前仓库源码看实际处理比文档表格更细DECIMAL 精度上限 38、默认精度 18、最大小数位 38、默认小数位 3见DuckDBTypeConverter.java#L80-L83。当精度或小数位超限时会做截断并输出 warning 日志小数位为负时归 0TIMESTAMP WITH TIME ZONE 映射为 OFFSET_DATE_TIME 类型见DuckDBTypeConverter.java#L169-L171即带时区偏移的时间类型而非普通 TIMESTAMP。该行为在 DuckDBTypeConverterTest 中有明确断言复杂类型 ARRAY / STRUCT / MAP 实际映射为 STRING默认长度 65535转换时会输出 warning 日志提示复杂类型已映射为 STRING可考虑使用 JSON 序列化见DuckDBTypeConverter.java#L176-L184而不是 BYTESBLOB 才映射为 BYTESPrimitiveByteArrayType无符号整数族UTINYINT、USMALLINT、UINTEGER、UBIGINT会分别落到对应的有符号 SeaTunnel 类型HUGEINT / UHUGEINT / BIGNUM 统一映射为DECIMAL(38,0)遇到未知类型如geography时会回退为 STRING 并输出 warningDuckDBTypeConverter.java#L185-L189。上述行为均以当前仓库源码为准若你使用的 SeaTunnel 发行版本不同映射结果可能略有差异建议以对应版本源码为准。并行读取Parallel Reader原理JDBC Source 支持对表数据进行并行读取。SeaTunnel 会按一定规则将表数据切分为多个 split再交给多个 reader 并行消费reader 数量由作业的parallelism决定。Split 键选择规则显式指定优先若配置了partition_column直接使用该列计算 split该列必须属于受支持的 split 数据类型自动推导兜底若未配置partition_columnSeaTunnel 会读取表 schema 获取主键Primary Key与唯一索引Unique Index。当主键/唯一索引包含多列时取其中第一个属于受支持 split 数据类型的列用于切分。例如表主键为(guid, name varchar)由于guid不属于受支持类型会自动改用name列切分。受支持的 split 数据类型String字符串Number数值类型int、bigint、decimal 等Date日期与切片相关的参数参数说明split.size单个 split 包含的行数表在读取时按此行数被切分为多个 splitpartition_column [string]用于切分数据的列名partition_upper_bound [BigDecimal]partition_column扫描最大值不设置时 SeaTunnel 查询数据库获取 max 值partition_lower_bound [BigDecimal]partition_column扫描最小值不设置时 SeaTunnel 查询数据库获取 min 值partition_num [int]需要切分成的 split 数量仅支持正整数默认等于作业并行度。官方不推荐使用正确做法是通过split.size控制切片数量注意partition_num在官方文档中标注为不推荐使用因为通过split.size控制每片行数更能适配数据量的动态变化。无法切分时的行为如果表既没有主键/唯一索引也未设置partition_column则该表以单并发single concurrency方式运行。对于这种场景可开启enable_concurrent_read false让源在快照阶段跳过切片分析、按单个 split 读取适合无索引的大表该选项定义于 JdbcSourceOptions。单表读取与多表读取的配置选择官方 Tips 给出两条核心建议单表读取优先用table_path替代query。配置table_path会自动开启自动切片auto split可通过split.*参数调整切片策略多表读取使用table_list。配置table_list同样会自动开启自动切片。从 DuckDBDialect 源码可见DuckDB 方言将默认数据库名设为default、默认 schema 设为main表标识符统一用双引号包裹如main.table1。因此table_path写作main.user_events等价于main.user_events。同时该方言刻意不支持 UPSERT 语义getUpsertStatement返回空Optional这是出于批量 ETL 负载与追加写入优化的设计取舍。实战任务示例以下示例均以本地 DuckDB 数据库文件/tmp/test.db为数据源Sink 使用 Console 插件将结果输出到控制台可直接替换数据库路径与表名后运行。示例一单并行简单查询该示例以单并行度查询测试库中的user_events表并输出全部字段你也可以通过修改query指定要查询的字段实现列裁剪。# Defining the runtime environment env { parallelism 4 job.mode BATCH } source{ Jdbc { url jdbc:duckdb:/tmp/test.db driver org.duckdb.DuckDBDriver connection_check_timeout_sec 100 username duckdb password query select * from user_events limit 16 } } transform { # 如需了解 transform 插件的更多配置请参考 SeaTunnel 官方 transform 文档 } sink { Console {} }示例二按 partition_column 并行读取通过partition_column指定切分列配合split.size控制每片行数。上下边界可通过partition_lower_bound/partition_upper_bound显式指定注释掉则自动获取。env { parallelism 4 job.mode BATCH } source { Jdbc { url jdbc:duckdb:/tmp/test.db driver org.duckdb.DuckDBDriver connection_check_timeout_sec 100 username duckdb password query select * from user_events partition_column id split.size 10000 # Read start boundary #partition_lower_bound ... # Read end boundary #partition_upper_bound ... } } sink { Console {} }示例三按主键或唯一索引自动并行配置table_path会开启自动切片读取时 SeaTunnel 会自动从表的主键/唯一索引中挑选合适的切分列。示例中同时保留了query用于在自动切片基础上限定读取范围。env { parallelism 4 job.mode BATCH } source { Jdbc { url jdbc:duckdb:/tmp/test.db driver org.duckdb.DuckDBDriver connection_check_timeout_sec 100 username password table_path main.user_events query select * from main.user_events split.size 10000 } } sink { Console {} }示例四指定并行边界显式声明partition_lower_bound与partition_upper_bound后只读取上下边界范围内的数据相比全表扫描更高效。同时可通过properties向 DuckDB 传递额外连接参数例如线程数与内存上限。source { Jdbc { url jdbc:duckdb:/tmp/test.db driver org.duckdb.DuckDBDriver connection_check_timeout_sec 100 username duckdb password # Define query logic as required query select * from user_events partition_column id # Read start boundary partition_lower_bound 1 # Read end boundary partition_upper_bound 500 partition_num 10 properties { threads4 memory_limit4GB } } }说明DuckDB 允许在 JDBC URL 或properties中传递连接参数如threads、memory_limit。官方文档指出当 properties 与 URL 参数同名时DuckDB 中 properties 优先于 URL。memory_limit用于限制 DuckDB 分析引擎可使用的内存上限threads控制其内部并行线程数合理设置有助于在受限环境中稳定运行。示例五多表读取通过table_list一次读取多张表。每张表可以独立配置query实现行/列过滤未配置query的表默认读取全量字段。注释部分展示了where_condition全局过滤与split.size切片行数的用法。env { job.mode BATCH parallelism 4 } source { Jdbc { url jdbc:duckdb:/tmp/test.db driver org.duckdb.DuckDBDriver connection_check_timeout_sec 100 username duckdb password table_list [ { table_path main.table1 }, { table_path main.table2 # Use query filter rows columns query select id, name from main.table2 where id 100 } ] #where_condition where id 100 #split.size 8096 } } sink { Console {} }源码级原理补充方言、URL 解析与 Catalog为了更透彻地理解该连接器可以关注以下四个源码切入点方言识别与注册DuckDBDialectFactory 通过AutoService(JdbcDialectFactory.class)自动注册acceptsURL以 URL 前缀jdbc:duckdb:判定归属DuckDBDialect 负责表路径解析、标识符引用双引号与类型映射器的装配URL 解析DuckDBURLParser 使用正则^jdbc:duckdb:(?path[^?]*?)(?suffix\?.*)?$提取数据库文件路径与查询参数后缀天然兼容jdbc:duckdb:memory:内存库形式host/port 记为 localhost/0Catalog 支持DuckDBCatalogFactory 提供表结构推断能力可选参数包括schema、decimal_type_narrowing、handle_blob_as_string默认 schema 为main这为table_path/table_list自动切片时的 schema 推断提供了基础端到端验证仓库内的 DuckDBSourceAndSinkTest 展示了覆盖 BOOLEAN、TINYINT、HUGEINT、无符号整数族、REAL、DECIMAL、VARCHAR/TEXT/CHAR/BPCHAR、BLOB、DATE/TIME/TIMESTAMP/TIMESTAMPTZ、INTERVAL、UUID 等全部类型的一张建表语句并通过真实 JDBC 连接完成 Source 读取与 Sink 写入的完整流程验证DuckDBConnectDryRunValidationTest 则验证了--dry-run connect钩子下的 schema 推断与连通性检查。如果你需要在批处理作业中快速验证连接配置可以参考仓库根目录下的 config/v2.batch.config.template 模板将其中 source 部分替换为上述 DuckDB 配置即可。常见注意点汇总无法切分则单并发表无主键/唯一索引且未设置partition_column时作业以单并发运行吞吐会受限建议为表补充主键或显式配置切分列优先table_path/table_list单表用table_path多表用table_list二者均自动开启切片query更灵活但需要自己控制切分边界驱动放置位置随引擎变化Spark/Flink 放plugins/Zeta 放lib/放错目录会报 ClassNotFound版本差异DuckDB 不同版本驱动类名可能不同务必核对实际驱动包类型映射以源码为准当前仓库中 TIMESTAMP WITH TIME ZONE 映射为带时区偏移的时间类型、ARRAY/STRUCT/MAP 映射为 STRING若与文档表格不一致以发行版本对应源码为准表路径唯一性table_list多表场景下各表table_path必须唯一否则作业校验失败。总结SeaTunnel JDBC DuckDB Source 是打通本地嵌入式 OLAP 分析库与分布式数据管道的轻量桥梁它没有网络服务依赖只需驱动 jar 与文件路径即可接入通过partition_column、主键/唯一索引自动切片与split.size等机制可以获得稳定的并行读取能力table_path/table_list让单表与多表批量读取的配置成本都保持在极低水平。结合本指南中的参数表、类型映射与五个可直接运行的示例你可以快速将 DuckDB 中的数据导入 SeaTunnel 支持的任意下游存储。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表