SeaTunnel Kingbase(人大金仓)JDBC Source 连接器:配置详解、类型映射与源码级拆分策略解析
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文基于 SeaTunnel 仓库中 Kingbase 源连接器官方文档(docs/en/connectors/source/Kingbase.md),系统讲解如何通过 JDBC 方式读取人大金仓 KingbaseES 数据库数据:涵盖驱动部署、完整 Source 配置参数、数据类型映射、并行分片(partition)与表路径(table_path)等实战配置,并结合仓库中connector-jdbc模块的 Kingbase 方言实现源码,深入剖析方言识别、类型转换 fallback 机制与大表切分策略,帮助你既能快速配置一个可运行的 Kingbase 同步任务,又能理解底层实现的选型依据。
一、连接器能力概览
Kingbase 是 SeaTunnel 通过 JDBC 方式支持的国产数据库连接器,KingbaseES V8R6(8.6)与 PostgreSQL 兼容,因此该连接器在connector-jdbc模块内以独立方言(Dialect)形式实现。根据官方文档与源码,其能力矩阵如下:
| 特性 | 支持情况 | 说明 |
|---|---|---|
| batch(批处理) | 支持 | 以job.mode = "BATCH"运行 |
| stream(流式) | 不支持 | 纯 JDBC 全量读取,无 CDC 能力 |
| exactly-once | 不支持 | 批模式下无源端精确一次语义 |
| column projection(列投影) | 支持 | 下游只取部分列时,查询自动裁剪列 |
| parallelism(并行度) | 支持 | 通过partition_*系列参数按分片并行读取 |
| support user-defined split(自定义拆分) | 支持 | 可指定分片边界控制读取范围 |
文档声明支持的引擎为Spark、Flink、SeaTunnel Zeta,即该 Source 配置可在这三类引擎中运行。
二、数据源信息与驱动部署
2.1 支持的数据源版本
Kingbase 连接器的接入要素如下:
| Datasource | 支持版本 | Driver 类名 | JDBC URL 示例 |
|---|---|---|---|
| Kingbase | 8.6 | com.kingbase8.Driver | jdbc:kingbase8://localhost:54321/db_test |
由于 Kingbase 官方驱动不在 SeaTunnel 发行包的依赖中,需要手动下载驱动 jar 并拷贝到插件目录:
# 下载 kingbase8-8.6.0.jar 后执行: cp kingbase8-8.6.0.jar $SEATUNNEL_HOME/plugins/jdbc/lib/2.2 方言如何被识别:源码证据
从源码结构看,JDBC 连接器通过 SPI 注册的方言工厂识别数据库。KingbaseDialectFactory.java 中的关键逻辑非常直接:
@Override public boolean acceptsURL(String url) { return url.startsWith("jdbc:kingbase8:"); } @Override public JdbcDialect create() { return new KingbaseDialect(); }也就是说,只要url以jdbc:kingbase8:开头,SeaTunnel 就会自动选用 Kingbase 方言(类型映射、标识符引用规则、行转换器等全部走 Kingbase 专属实现)。这解释了为什么driver必须配置为com.kingbase8.Driver而 URL 前缀不可写成jdbc:postgresql:——后者会命中 PostgreSQL 方言。
三、数据类型映射
3.1 官方映射表
Kingbase 到 SeaTunnel 的类型映射如下(摘自官方文档):
| Kingbase 数据类型 | SeaTunnel 数据类型 |
|---|---|
| BOOL | BOOLEAN |
| INT2 | SHORT |
| SMALLSERIAL / SERIAL / INT4 | INT |
| INT8 / BIGSERIAL | BIGINT |
| FLOAT4 | FLOAT |
| FLOAT8 | DOUBLE |
| NUMERIC | DECIMAL(precision, scale),取列定义的小数位数 |
| BPCHAR / CHARACTER / VARCHAR / TEXT | STRING |
| TIMESTAMP | LOCALDATETIME |
| TIME | LOCALTIME |
| DATE | LOCALDATE |
| 其他类型 | 暂不支持 |
3.2 源码级深挖:fallback 类型转换
上表是文档口径的基础映射,而仓库源码中的 KingbaseTypeConverter.java 揭示了更完整的实际行为:
@AutoService(TypeConverter.class) public class KingbaseTypeConverter extends PostgresTypeConverter {Kingbase 类型转换器继承自PostgresTypeConverter(印证了 KingbaseES 对 PostgreSQL 的兼容性),并在此之上做三层扩展:
- Kingbase 独有类型的兜底处理:
convert()先调用父类PostgresTypeConverter转换,父类抛异常时进入 switch 兜底分支,支持:TINYINT→BYTE_TYPE(PG 无此类型,Kingbase 兼容 MySQL 语法提供);MONEY→DECIMAL(38, 18);BLOB→ 二进制类型(PrimitiveByteArrayType,长度上限 1GB);CLOB→STRING(长度上限 1GB);BIT(M)→ 二进制类型,按M/8向上取整换算成字节长度。
- 跨库兼容类型:由于 KingbaseES 具备 MySQL/Oracle 兼容模式,转换器还纳入了
INT、MEDIUMINT、DATETIME、BLOB/TEXT系列(MySQL 风格)、NUMBER、VARCHAR2、NVARCHAR2、XML(Oracle 风格)、DATETIME2、DATETIMEOFFSET(SQL Server 风格)等类型名,避免兼容模式建表时报"不支持的类型"错误。 - 超限时降级而非报错:例如
DATETIME的 scale 超过MAX_TIMESTAMP_SCALE时会截断到最大 scale 并打 warn 日志,而不中断任务。
读取路径上,KingbaseTypeMapper.java 从ResultSetMetaData中取列名、getColumnTypeName()原始类型、精度与小数位,组装成BasicTypeDefine后交给KingbaseTypeConverter完成最终映射——这就是文档中NUMERIC能自动映射出精确DECIMAL(precision, scale)的实现来源。行数据的转换则由 KingbaseJdbcRowConverter.java 负责,与 PostgreSQL 的 JDBC 行转换逻辑保持一致。
四、Source Options 完整参数表
以下参数与官方文档完全对齐(Kingbase 与 PostgreSQL 等其他 JDBC 方言共用这套 Source 选项,定义见 JdbcSourceOptions.java):
| Name | Type | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL,示例:jdbc:kingbase8://localhost:54321/test |
| driver | String | 是 | - | 连接远程数据源的 JDBC 类名,应为com.kingbase8.Driver |
| username | String | 否 | - | 连接实例用户名;旧配置键user仍作为 fallback 被接受 |
| password | String | 否 | - | 连接实例密码 |
| query | String | 是 | - | 查询语句 |
| connection_check_timeout_sec | Int | 否 | 30 | 校验连接可用性时等待数据库操作完成的秒数 |
| partition_column | String | 否 | - | 并行分片列,仅支持数值类型和字符串类型列 |
| partition_lower_bound | BigDecimal | 否 | - | 分片列最小扫描值,不设置时 SeaTunnel 会查询数据库获取 min 值 |
| partition_upper_bound | BigDecimal | 否 | - | 分片列最大扫描值,不设置时 SeaTunnel 会查询数据库获取 max 值 |
| partition_num | Int | 否 | 作业并行度 | 分片数量,仅支持正整数,默认为作业并行度 |
| fetch_size | Int | 否 | 0 | 大结果集查询的行抓取大小,减少数据库往返次数以提升性能;0 表示使用 JDBC 默认值 |
| use_regex | Boolean | 否 | false | 控制table_path是否按正则匹配。为true时按正则模式匹配,否则按精确路径 |
| table_path | String | 否 | - | 表全路径,可替代query,例如test_schema.table1 |
| table_list | Array | 否 | - | 要读取的表列表,可替代table_path,例如[{ table_path = "testdb.table1"}, {table_path = "testdb.table2", query = "select id, name from testdb.table2"}] |
| where_condition | String | 否 | - | 作用于所有表/查询的公共行过滤条件,必须以where开头,例如where id > 100 |
| split.size | Int | 否 | 8096 | 按表读取时的单个 split 行数,表会被拆分为多个 split |
| split.even-distribution.factor.lower-bound | Double | 否 | 0.05 | 分片列分布因子的下界。分布因子 = (MAX(id) - MIN(id) + 1) / 行数,落入 [下界, 上界] 区间时按均匀分布优化切分;低于下界则视为不均匀分布,在估计分片数超过split.sample-sharding.threshold时改用采样切分策略 |
| split.even-distribution.factor.upper-bound | Double | 否 | 100 | 分布因子上界,语义同上,超出上界同样触发采样切分评估 |
| split.sample-sharding.threshold | Int | 否 | 1000 | 触发采样切分策略的估计分片数阈值(估算行数 / split.size)。超出阈值时启用采样切分以高效处理大表。注意:当前仓库源码中该选项默认值为 1000(见 JdbcSourceOptions.java 的defaultValue(1000)),与文档表格中 10000 的写法不一致,建议以源码为准 |
| split.inverse-sampling.rate | Int | 否 | 1000 | 采样切分策略的采样率倒数,1000 表示 1/1000 采样率;数值越大采样越稀疏,适合超大表 |
| common-options | - | 否 | - | Source 插件公共参数,见 Source Common Options 文档(docs/en/connectors/common-options/source-common-options.md) |
4.1 分片(Split)策略的源码印证
上述split.*参数并非空配置,它们在源码中有明确的实现载体:
- 均匀分布判定与采样切分由 DynamicChunkSplitter.java 实现,源码注释明确说明
maxChunkCount由split.sample-sharding.threshold提供,分布因子超出上下界且估计分片数超阈值时切换到基于split.inverse-sampling.rate的采样切分; - 所有默认值(
split.size = 8096、分布因子0.05/100、采样阈值1000、采样率倒数1000)均可在 JdbcSourceOptions.java 中以Options.key(...).defaultValue(...)形式逐项核对。
4.2 Tips:并行与单并发的选择
若未设置
partition_column,任务将以单并发运行;设置了partition_column后,会按任务并行度并行执行。
五、实战配置示例
5.1 简单全量读取(Simple)
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 插件配置方式,可查阅官方 transforms/sql 文档 } sink { Console {} }5.2 并行读取整表(Parallel)
通过配置分片字段实现并行读取,适合整表抽取场景。
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 } }5.3 指定上下边界的并行读取(Parallel Boundary)
显式指定分片列的上下界可以省去 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 } }5.4 按 Schema 限定表名(Query With Schema Name)
Kingbase 表名通常写作
schema.table。连接账号既可以用username,也可以使用旧配置键user(作为 fallback 被接受)。
source { Jdbc { driver = "com.kingbase8.Driver" url = "jdbc:kingbase8://localhost:54321/test" user = "SYSTEM" password = "123456" query = "select * from public.e2e_table_source" } }5.5 仓库内置 E2E 任务配置参考
仓库中 Kingbase 的端到端集成测试提供了可直接参考的完整任务配置(Source 读public.e2e_table_source,Sink 批量写入public.e2e_table_sink):jdbc_kingbase_source_and_sink.conf:
env { parallelism = 1 job.mode = "BATCH" } source { jdbc { driver = "com.kingbase8.Driver" url = "jdbc:kingbase8://e2e_KINGBASEDb:54321/test" user = "SYSTEM" password = "123456" query = "select * from public.e2e_table_source" } } sink { jdbc { driver = "com.kingbase8.Driver" url = "jdbc:kingbase8://e2e_KINGBASEDb:54321/test" user = "SYSTEM" password = "123456" query = "INSERT INTO public.e2e_table_sink (c1, c2, c3, ...) VALUES (?, ?, ?, ...)" } }对应的测试入口为 JdbcKingbaseIT.java,它基于 Testcontainers 启动 Kingbase 容器并校验 Source 与 Sink 两侧的列值一致性,可作为验证连接器可用性的参照。
六、Kingbase 方言的更多实现细节
除了读取侧,仓库中的 Kingbase 方言实现还覆盖了标识符引用、upsert 与建表等能力,理解这些细节有助于排查 DDL/标识符相关问题:
- 标识符引用规则:KingbaseDialect.java 的
quoteIdentifier()对含.的多段名称逐段加双引号(如"schema"."table"),与 PostgreSQL 风格一致,避免大小写敏感与保留字问题;tableIdentifier()对库名单独加引号以解决 PG 系数据库名大小写不敏感的历史问题。方言还支持fieldIde字段命名风格配置(默认 ORIGINAL)。 - Upsert 语句:
getUpsertStatement()基于 PostgreSQL 兼容语法生成INSERT ... ON CONFLICT (pk) DO UPDATE SET col = EXCLUDED.col ...,写侧连接器可直接利用该主键冲突更新能力。 - 建表与表选项校验:KingbaseCreateTableSqlBuilder.java 与 KingbaseCatalog.java 实现了 catalog 能力;方言仅开放
tablespace与fillfactor两个建表选项,其中fillfactor会被校验为 10–100 的整数、tablespace会拒绝含引号、换行或分号的值,防止 DDL 注入——从源码结构看,这些校验失败会抛出CONFIG_VALIDATION_FAILED错误码的配置校验异常。 - 变更历史:该连接器自 2.3.4 版本(提交 "jdbc connector supports Kingbase database (#4803)")加入,后续随 JDBC 连接器统一演进,完整变更记录见 connector-jdbc Changelog。
七、小结与适用前提
- 适用前提:KingbaseES 8.6 数据源、驱动 jar 已放入
$SEATUNNEL_HOME/plugins/jdbc/lib/、JDBC URL 以jdbc:kingbase8:开头;引擎可选 Zeta / Flink / Spark,运行模式为 BATCH。 - 整表抽取建议优先使用
table_path/table_list搭配split.*参数让连接器自动按split.size拆片;已有明确主键/自增列且希望控制并发粒度时,再使用partition_column系列参数并显式给定partition_lower_bound/partition_upper_bound省去 min/max 探测。 - 类型映射上,文档表格之外的
TINYINT、MONEY、BLOB、CLOB、BIT以及 MySQL/Oracle 兼容模式类型名已被 KingbaseTypeConverter.java 显式兜底,遇到文档未列类型时建议先核对该文件的 switch 分支,确认是否已在当前版本支持。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考