SeaTunnel Kingbase 源连接器实战指南:JDBC 读取、并行分片与类型映射全解析
2026/9/19 2:59:50 网站建设 项目流程

SeaTunnel Kingbase 源连接器实战指南:JDBC 读取、并行分片与类型映射全解析

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

Kingbase(人大金仓)是国内广泛使用的 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

该连接器支持以下三类运行引擎:

  • Spark
  • Flink
  • SeaTunnel 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 数据类型
BOOLBOOLEAN
INT2SHORT
SMALLSERIAL
SERIAL
INT4
INT
INT8
BIGSERIAL
BIGINT
FLOAT4FLOAT
FLOAT8DOUBLE
NUMERICDECIMAL
BPCHAR
CHARACTER
VARCHAR
TEXT
STRING
TIMESTAMPLOCALDATETIME
TIMELOCALTIME
DATELOCALDATE
其他数据类型暂不支持

类型转换的源码实现

上述映射关系在 KingbaseTypeConverter 中实现。值得关注的是,Kingbase 作为一款兼容多数据库语法的国产数据库,其类型转换器继承自PostgresTypeConverter,并在此基础上补充了 MySQL、Oracle、SQL Server 兼容类型以及 Kingbase 特有类型的转换分支,包括:

  • Kingbase 特有类型TINYINTBYTEMONEYDECIMAL(38,18)BLOBBYTESCLOBSTRINGBIT(M)BYTES(按位折算字节长度);
  • MySQL 兼容类型:如MEDIUMINTINTYEAR等映射为INTDATETIME映射为LOCAL_DATETIMEBINARY/VARBINARY及各类BLOB映射为BYTESTINYTEXT/MEDIUMTEXT/LONGTEXT映射为STRING
  • Oracle 兼容类型NUMBER按精度/小数位映射为DECIMALVARCHAR2NVARCHAR2NCHARLONGROWIDCLOB等映射为STRING
  • SQL Server 兼容类型DATETIME2SMALLDATETIME映射为LOCAL_DATETIMEDATETIMEOFFSET映射为OFFSET_DATE_TIME

对应的单元测试 KingbaseTypeConverterTest 覆盖了boolBOOLEANint2SHORTint4INTint8LONG、浮点类型以及不支持类型抛异常等场景,可作为确认映射行为的第一手依据。

从源码结构看,如果 Kingbase 表中存在上表之外的类型(例如数组、JSON 等复杂类型),类型转换会抛出SeaTunnelRuntimeException(提示"暂不支持"),此时需要在上游 SQL 中通过显式CAST将列转换为支持的类型后再同步。

源选项(Source Options)完整说明

参数名类型必须默认值描述
urlString-JDBC 连接 URL。参考示例:jdbc:kingbase8://localhost:54321/test
driverString-连接远程数据源的 JDBC 驱动类名,应为com.kingbase8.Driver
usernameString-连接实例用户名。旧配置名user仍可作为兼容写法使用
passwordString-连接实例密码
queryString-查询语句
connection_check_timeout_secInt30等待用于验证连接的数据库操作完成的时间(秒)
partition_columnString-用于并行性分割的列名,仅支持数值类型列和字符串类型列
partition_lower_boundBigDecimal-partition_column 的最小值用于扫描,如果未设置,SeaTunnel 将查询数据库获取最小值
partition_upper_boundBigDecimal-partition_column 的最大值用于扫描,如果未设置,SeaTunnel 将查询数据库获取最大值
partition_numIntjob parallelism分割数量,仅支持正整数。默认值是任务并行度
fetch_sizeInt0对于返回大量对象的查询,可配置查询中使用的行提取大小,通过减少满足选择条件所需的数据库命中次数来提高性能。零表示使用 JDBC 默认值
use_regexBooleanfalse控制表路径的正则表达式匹配。设为true时,table_path将被视为正则表达式模式;设为false或未指定时,table_path被视为精确路径(不进行正则匹配)
table_pathString-表的完整路径,可用此配置代替query。示例:"testdb.table1"
table_listArray-要读取的表的列表,可用此配置代替table_path。示例:[{ table_path = "testdb.table1"}, {table_path = "testdb.table2", query = "select * id, name from testdb.table2"}]
where_conditionString-所有表/查询的通用行过滤条件,必须以where开头。例如where id > 100
split.sizeInt8096表的分割大小(行数),读取表时,捕获的表会被分割成多个分片
split.even-distribution.factor.lower-boundDouble0.05分片键分布因子的下限。该因子用于判断表数据分布是否均匀。若计算得到的分布因子大于等于该下限(即(MAX(id) - MIN(id) + 1) / 行数),则会对表的分片进行优化以确保数据均匀分布;反之,若分布因子较低,则表数据被视为分布不均匀。若估算的分片数量超过sample-sharding.threshold指定的值,则采用基于采样的分片策略
split.even-distribution.factor.upper-boundDouble100分片键分布因子的上限。若计算得到的分布因子小于等于该上限,则对表的分片进行均匀分布优化;反之,若分布因子较大,则表数据被视为分布不均匀,且当估算分片数超过sample-sharding.threshold时采用基于采样的分片策略
split.sample-sharding.thresholdInt1000触发样本分片策略的估算分片数阈值。当分布因子超出上下限范围,且估算分片数(大致行数 / 分片大小)超过此阈值时,将使用样本分片策略,有助于更高效地处理大型数据集
split.inverse-sampling.rateInt1000样本分片策略中使用的采样率的倒数。例如设置为 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 = 1partition_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)。userusername的兼容旧写法,二者等价。

按表路径读取(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_conditionwhere开头,作用于所有表/查询,可作为通用的行级过滤手段。

底层实现:方言、行转换与 Catalog

方言(Dialect)与标识符引用

KingbaseDialect 实现了JdbcDialect接口,方言名取自DatabaseIdentifier.KINGBASE。其实现细节包括:

  • 标识符引用quoteIdentifier采用双引号包裹标识符(如"column"),并对含.的多段标识符逐段加引号,符合 Kingbase/PostgreSQL 的大小写敏感语义;
  • Upsert 支持getUpsertStatement生成INSERT ... ON CONFLICT (pk) DO UPDATE SET col = EXCLUDED.col语法,说明该方言在 Sink 场景下也支持基于主键冲突的写模式;
  • 表选项校验validateTableOptions仅接受tablespacefillfactor两个表选项,其中fillfactor必须是 10~100 的整数,tablespace值不允许包含引号、换行与分号等非法字符,防止 DDL 注入。

行转换与类型映射

KingbaseDialect.getRowConverter()返回KingbaseJdbcRowConvertergetJdbcDialectTypeMapper()返回KingbaseTypeMapper,分别负责 JDBC 结果集到 SeaTunnel 行对象、以及 JDBC 元数据类型到 SeaTunnel 类型的双向转换。

Catalog 元数据读取

KingbaseCatalog 继承自AbstractJdbcCatalog,通过查询sys_classsys_namespacesys_attribute等系统表获取列名、类型、长度、精度、默认值与注释等元数据。它默认排除INFORMATION_SCHEMASYSAUDITSYSLOGICALSYS_CATALOGSYS_HMXLOG_RECORD_READ等系统 schema,避免在元数据遍历时污染业务表集合。配合 KingbaseCatalogFactory 与 KingbaseCreateTableSqlBuilder,该连接器在 Sink 场景下还能基于 CatalogTable 自动建表。其建表/类型行为由 KingbaseCatalogTest 等测试用例持续验证。

常见问题与调优建议

  1. 驱动未找到:确保kingbase8-8.6.0.jar已复制到$SEATUNNEL_HOME/plugins/jdbc/lib/,且版本与 Kingbase 服务端 8.6 匹配;连接串必须以jdbc:kingbase8://开头,否则方言工厂无法识别。
  2. 并行不生效:未设置partition_column时任务只能单并发运行;设置了partition_column但未设置partition_num时,分片数默认取环境并行度,可在env中调整parallelism
  3. 分片键选择partition_column仅支持数值类型列与字符串类型列;建议优先选择主键或高基数、分布均匀的列,以获得均衡的分片区间。
  4. 大表分片策略:当数据分布不均匀且估算分片数超过split.sample-sharding.threshold时,连接器会自动切换到基于采样的分片策略,可通过split.inverse-sampling.rate调整采样粒度。
  5. 类型不支持:遇到映射表中"其他数据类型",请在查询 SQL 中显式CAST为目标类型,例如将复杂类型转换为字符串或数值。
  6. 吞吐优化:对返回大量行的查询可设置fetch_size(如 1000~5000),通过减少数据库往返次数提升读取性能;同时建议开启并行分片并结合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),仅供参考

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

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

立即咨询