SeaTunnel Kingbase 源连接器实战指南:JDBC 读取、并行分片与类型映射全解析
发布时间:2026/9/19 2:59:54 作者:尧图编辑部 阅读量:1,286

SeaTunnel Kingbase 源连接器实战指南JDBC 读取、并行分片与类型映射全解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelKingbase人大金仓是国内广泛使用的 PostgreSQL 系国产关系型数据库本指南围绕 SeaTunnel 官方提供的JDBC Kingbase 源连接器展开讲解如何通过 HOCON 配置将 Kingbase 8.6 中的数据以批模式、并行分片方式高效读入 SeaTunnel 管道并深入解析其底层方言实现、数据类型映射与分片策略原理。读完本文你将掌握 Kingbase 作为数据源接入 SeaTunnel 的完整配置方法、并行优化手段以及连接器内部的驱动匹配、类型转换与 Catalog 元数据读取机制。本文主体基于 docs/zh/connectors/source/Kingbase.md并辅以 connector-jdbc 模块中 Kingbase 方言的真实源码与单元测试进行佐证。连接器概览与引擎支持Kingbase 源连接器通过 JDBC 读取外部数据源数据属于 connector-v2 体系中 Jdbc 连接器家族的一员。当前文档标记的支持连接器版本为 8.6对应的官方驱动为com.kingbase8.Driver。该连接器支持以下三类运行引擎SparkFlinkSeaTunnel Zeta在关键特性方面特性定义可参考 连接器 v2 特性说明当前版本的能力矩阵如下批处理BATCH✅ 支持流处理STREAM❌ 不支持精确一次Exactly Once❌ 不支持列投影Column Projection✅ 支持并行性Parallelism✅ 支持用户自定义 Split✅ 支持这意味着 Kingbase 源主要用于批式数据同步场景且天然具备并行分片读取能力适合大表全量抽取。支持的数据源与驱动依赖数据源信息数据源支持的版本驱动连接串示例Maven 坐标Kingbase8.6com.kingbase8.Driverjdbc:kingbase8://localhost:54321/db_testcn.com.kingbase:kingbase8:8.6.0kingbase8-8.6.0.jar数据库驱动部署请下载对应版本如kingbase8-8.6.0.jar并将其复制到$SEATUNNEL_HOME/plugins/jdbc/lib/工作目录。例如cp kingbase8-8.6.0.jar $SEATUNNEL_HOME/plugins/jdbc/lib/驱动放置完成后SeaTunnel 启动时即可在 JDBC 连接器的插件类加载路径中找到com.kingbase8.Driver。驱动匹配的源码依据从源码可以确认SeaTunnel 通过 SPI 机制AutoService(JdbcDialectFactory.class)注册 Kingbase 方言工厂。在 KingbaseDialectFactory 中acceptsURL方法通过前缀jdbc:kingbase8:识别连接串并据此创建对应的方言实例Override public boolean acceptsURL(String url) { return url.startsWith(jdbc:kingbase8:); }因此配置中url必须以jdbc:kingbase8://开头否则不会被识别为 Kingbase 方言。Kingbase 数据类型映射Kingbase 与 SeaTunnel 内部类型系统之间的映射关系如下表Kingbase 数据类型SeaTunnel 数据类型BOOLBOOLEANINT2SHORTSMALLSERIALSERIALINT4INTINT8BIGSERIALBIGINTFLOAT4FLOATFLOAT8DOUBLENUMERICDECIMALBPCHARCHARACTERVARCHARTEXTSTRINGTIMESTAMPLOCALDATETIMETIMELOCALTIMEDATELOCALDATE其他数据类型暂不支持类型转换的源码实现上述映射关系在 KingbaseTypeConverter 中实现。值得关注的是Kingbase 作为一款兼容多数据库语法的国产数据库其类型转换器继承自PostgresTypeConverter并在此基础上补充了 MySQL、Oracle、SQL Server 兼容类型以及 Kingbase 特有类型的转换分支包括Kingbase 特有类型TINYINT→BYTE、MONEY→DECIMAL(38,18)、BLOB→BYTES、CLOB→STRING、BIT(M)→BYTES按位折算字节长度MySQL 兼容类型如MEDIUMINT、INT、YEAR等映射为INTDATETIME映射为LOCAL_DATETIMEBINARY/VARBINARY及各类BLOB映射为BYTESTINYTEXT/MEDIUMTEXT/LONGTEXT映射为STRINGOracle 兼容类型NUMBER按精度/小数位映射为DECIMALVARCHAR2、NVARCHAR2、NCHAR、LONG、ROWID、CLOB等映射为STRINGSQL Server 兼容类型DATETIME2、SMALLDATETIME映射为LOCAL_DATETIMEDATETIMEOFFSET映射为OFFSET_DATE_TIME。对应的单元测试 KingbaseTypeConverterTest 覆盖了bool→BOOLEAN、int2→SHORT、int4→INT、int8→LONG、浮点类型以及不支持类型抛异常等场景可作为确认映射行为的第一手依据。从源码结构看如果 Kingbase 表中存在上表之外的类型例如数组、JSON 等复杂类型类型转换会抛出SeaTunnelRuntimeException提示暂不支持此时需要在上游 SQL 中通过显式CAST将列转换为支持的类型后再同步。源选项Source Options完整说明参数名类型必须默认值描述urlString是-JDBC 连接 URL。参考示例jdbc:kingbase8://localhost:54321/testdriverString是-连接远程数据源的 JDBC 驱动类名应为com.kingbase8.DriverusernameString否-连接实例用户名。旧配置名user仍可作为兼容写法使用passwordString否-连接实例密码queryString是-查询语句connection_check_timeout_secInt否30等待用于验证连接的数据库操作完成的时间秒partition_columnString否-用于并行性分割的列名仅支持数值类型列和字符串类型列partition_lower_boundBigDecimal否-partition_column 的最小值用于扫描如果未设置SeaTunnel 将查询数据库获取最小值partition_upper_boundBigDecimal否-partition_column 的最大值用于扫描如果未设置SeaTunnel 将查询数据库获取最大值partition_numInt否job parallelism分割数量仅支持正整数。默认值是任务并行度fetch_sizeInt否0对于返回大量对象的查询可配置查询中使用的行提取大小通过减少满足选择条件所需的数据库命中次数来提高性能。零表示使用 JDBC 默认值use_regexBoolean否false控制表路径的正则表达式匹配。设为true时table_path将被视为正则表达式模式设为false或未指定时table_path被视为精确路径不进行正则匹配table_pathString否-表的完整路径可用此配置代替query。示例testdb.table1table_listArray否-要读取的表的列表可用此配置代替table_path。示例[{ table_path testdb.table1}, {table_path testdb.table2, query select * id, name from testdb.table2}]where_conditionString否-所有表/查询的通用行过滤条件必须以where开头。例如where id 100split.sizeInt否8096表的分割大小行数读取表时捕获的表会被分割成多个分片split.even-distribution.factor.lower-boundDouble否0.05分片键分布因子的下限。该因子用于判断表数据分布是否均匀。若计算得到的分布因子大于等于该下限即(MAX(id) - MIN(id) 1) / 行数则会对表的分片进行优化以确保数据均匀分布反之若分布因子较低则表数据被视为分布不均匀。若估算的分片数量超过sample-sharding.threshold指定的值则采用基于采样的分片策略split.even-distribution.factor.upper-boundDouble否100分片键分布因子的上限。若计算得到的分布因子小于等于该上限则对表的分片进行均匀分布优化反之若分布因子较大则表数据被视为分布不均匀且当估算分片数超过sample-sharding.threshold时采用基于采样的分片策略split.sample-sharding.thresholdInt否1000触发样本分片策略的估算分片数阈值。当分布因子超出上下限范围且估算分片数大致行数 / 分片大小超过此阈值时将使用样本分片策略有助于更高效地处理大型数据集split.inverse-sampling.rateInt否1000样本分片策略中使用的采样率的倒数。例如设置为 1000表示采样过程中应用 1/1000 的采样率。该选项可控制采样粒度、影响最终分片数量特别适用于数据量极大、通常需要较低采样率的场景common-options否-源插件通用参数详见 源通用选项参数实现的源码佐证以上核心参数在 JdbcSourceOptions 中以Option定义其中几个值得注意的实现细节split.even-distribution.factor.upper-bound默认值为100.0dlower-bound默认值为0.05d分布因子通过(MAX(id) - MIN(id) 1) / rowCount计算用于判断数据分布是否均匀split.sample-sharding.threshold默认值为 1000 个分片注意文档参数表正文写为 10000而源码注释与默认值均为 1000实际生效值以源码 JdbcSourceOptions#L78-L88 为准除文档列出的参数外源码还提供了split.allow-sampling默认true关闭后回退到非均匀迭代式分片、use_select_count是否用select count统计行数、skip_analyze跳过表行数分析、enable_concurrent_read默认true快照阶段并发读取等可选参数可按需在配置中启用partition_num默认取任务并行度即env.parallelism与文档一致。使用提示如果未设置partition_column任务将以单并发运行如果设置了partition_column任务将根据配置的并发度partition_num或环境并行度并行执行。任务示例以下示例均基于 HOCON 配置语法可直接放入 SeaTunnel 配置文件运行。完整的env/source/transform/sink骨架可参考仓库根目录的 v2.batch.config.template。简单示例单并发全量查询env { parallelism 2 job.mode BATCH } source { Jdbc { driver com.kingbase8.Driver url jdbc:kingbase8://localhost:54321/db_test username root password query select * from source } } transform { # 此处可按需配置 transform 插件 } sink { Console {} }该示例未设置partition_column因此整个查询以单任务读取job.mode必须为BATCH当前连接器不支持流模式。Sink 使用Console便于在控制台直接观察读取结果。并行示例按分片字段并行读取使用您配置的分片字段和分片数据并行读取查询表。如果您想读取整个表可以这样做source { Jdbc { driver com.kingbase8.Driver url jdbc:kingbase8://localhost:54321/db_test username root password query select * from source # 并行分片读取字段 partition_column id # 分片数量 partition_num 10 } }设置partition_column id后SeaTunnel 会基于id列将查询拆分为 10 个范围分片并行读取。若partition_lower_bound/partition_upper_bound未指定连接器会先向数据库查询该列的MIN/MAX值作为边界。并行边界示例显式指定上下界根据您配置的上下边界读取数据源更高效。source { Jdbc { driver com.kingbase8.Driver url jdbc:kingbase8://localhost:54321/db_test username root password query select * from source partition_column id partition_num 10 # 读取开始边界 partition_lower_bound 1 # 读取结束边界 partition_upper_bound 500 } }显式给出partition_lower_bound 1与partition_upper_bound 500可以省去连接器查询MIN/MAX的额外开销同时把分片范围精确限定在[1, 500]在大表上能显著提升分片划分效率。使用 Schema 表名查询Kingbase 表名通常写成schema.table。连接用户名可以使用username也可以使用兼容写法user。source { Jdbc { driver com.kingbase8.Driver url jdbc:kingbase8://localhost:54321/test user SYSTEM password 123456 query select * from public.e2e_table_source } }Kingbase 沿用了 PostgreSQL 的 schema 组织方式因此查询语句中建议显式写出schema.table如public.e2e_table_source。user是username的兼容旧写法二者等价。按表路径读取table_path / table_list除了query还可以直接用表路径驱动读取。例如source { Jdbc { driver com.kingbase8.Driver url jdbc:kingbase8://localhost:54321/db_test username root password table_path public.source_table } }多表场景则使用table_list并可为每张表单独指定查询与过滤条件source { Jdbc { driver com.kingbase8.Driver url jdbc:kingbase8://localhost:54321/db_test username root password table_list [ { table_path public.table1 }, { table_path public.table2, query select id, name from public.table2 } ] where_condition where id 100 } }where_condition以where开头作用于所有表/查询可作为通用的行级过滤手段。底层实现方言、行转换与 Catalog方言Dialect与标识符引用KingbaseDialect 实现了JdbcDialect接口方言名取自DatabaseIdentifier.KINGBASE。其实现细节包括标识符引用quoteIdentifier采用双引号包裹标识符如column并对含.的多段标识符逐段加引号符合 Kingbase/PostgreSQL 的大小写敏感语义Upsert 支持getUpsertStatement生成INSERT ... ON CONFLICT (pk) DO UPDATE SET col EXCLUDED.col语法说明该方言在 Sink 场景下也支持基于主键冲突的写模式表选项校验validateTableOptions仅接受tablespace与fillfactor两个表选项其中fillfactor必须是 10100 的整数tablespace值不允许包含引号、换行与分号等非法字符防止 DDL 注入。行转换与类型映射KingbaseDialect.getRowConverter()返回KingbaseJdbcRowConvertergetJdbcDialectTypeMapper()返回KingbaseTypeMapper分别负责 JDBC 结果集到 SeaTunnel 行对象、以及 JDBC 元数据类型到 SeaTunnel 类型的双向转换。Catalog 元数据读取KingbaseCatalog 继承自AbstractJdbcCatalog通过查询sys_class、sys_namespace、sys_attribute等系统表获取列名、类型、长度、精度、默认值与注释等元数据。它默认排除INFORMATION_SCHEMA、SYSAUDIT、SYSLOGICAL、SYS_CATALOG、SYS_HM、XLOG_RECORD_READ等系统 schema避免在元数据遍历时污染业务表集合。配合 KingbaseCatalogFactory 与 KingbaseCreateTableSqlBuilder该连接器在 Sink 场景下还能基于 CatalogTable 自动建表。其建表/类型行为由 KingbaseCatalogTest 等测试用例持续验证。常见问题与调优建议驱动未找到确保kingbase8-8.6.0.jar已复制到$SEATUNNEL_HOME/plugins/jdbc/lib/且版本与 Kingbase 服务端 8.6 匹配连接串必须以jdbc:kingbase8://开头否则方言工厂无法识别。并行不生效未设置partition_column时任务只能单并发运行设置了partition_column但未设置partition_num时分片数默认取环境并行度可在env中调整parallelism。分片键选择partition_column仅支持数值类型列与字符串类型列建议优先选择主键或高基数、分布均匀的列以获得均衡的分片区间。大表分片策略当数据分布不均匀且估算分片数超过split.sample-sharding.threshold时连接器会自动切换到基于采样的分片策略可通过split.inverse-sampling.rate调整采样粒度。类型不支持遇到映射表中其他数据类型请在查询 SQL 中显式CAST为目标类型例如将复杂类型转换为字符串或数值。吞吐优化对返回大量行的查询可设置fetch_size如 10005000通过减少数据库往返次数提升读取性能同时建议开启并行分片并结合partition_lower_bound/partition_upper_bound缩小扫描范围。变更日志该连接器的历史变更记录可在 connector-jdbc 变更日志 中查阅原文档通过ChangeLog /组件内嵌该文件其中记录了 JDBC 连接器家族各版本的修复与增强可作为升级评估的参考依据。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考