SeaTunnel Kingbase Sink 连接器完全指南:JDBC 配置、类型映射与实战写入
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文围绕 Kingbase Sink 连接器文档 展开,系统讲解 SeaTunnel 中面向人大金仓 Kingbase 数据库(8.6)的 JDBC Sink 接入方式。你将掌握驱动依赖的安装、JDBC 连接参数与全部 Sink 选项、Kingbase 与 SeaTunnel 数据类型的映射规则,以及三种可直接复用的任务配置(自定义 SQL 写入、自动生成 Sink SQL、带 Schema 表写入),并了解其底层 Dialect 实现原理。
连接器概览
Kingbase Sink 是 SeaTunnel 的 JDBC 系列 Sink 连接器之一,用于将上游数据批量写入人大金仓 Kingbase 数据库。它复用 connector-jdbc 统一的 JDBC 写入框架,通过独立的 Dialect 适配层(源码位于 internal/dialect/kingbase)实现 Kingbase 专属的 SQL 生成、标识符引用与类型转换逻辑。
支持连接器版本
- Kingbase 8.6
支持引擎
- Spark
- Flink
- SeaTunnel Zeta
关键特性
| 特性 | 支持情况 |
|---|---|
| 精确一次(Exactly Once) | ❌ 不支持 |
| CDC | ❌ 不支持 |
| 定时刷新 | ✅ 支持 |
需要特别说明的是:官方文档明确指出,连接器依赖XA 事务来保证精确一次语义,因此只有支持XA 事务的数据库才能启用精确一次(通过is_exactly_once=true),而Kingbase 目前不支持 XA 事务,所以该特性在 Kingbase 上不可用。与此对应的xa_data_source_class_name选项同样不适用于 Kingbase。
支持的数据源信息
| 数据源 | 支持的版本 | 驱动 | URL | Maven |
|---|---|---|---|---|
| Kingbase | 8.6 | com.kingbase8.Driver | jdbc:kingbase8://localhost:54321/db_test | kingbase8-8.6.0.jar |
从仓库的 connector-jdbc/pom.xml 可以看到,JDBC 连接器模块已将cn.com.kingbase:kingbase8以8.6.0版本声明为依赖(<kingbase8.version>8.6.0</kingbase8.version>),与实际支持的版本一致。
数据库依赖安装
使用该连接器前,需要下载对应的 JDBC 驱动 jar 并放入 SeaTunnel 的插件目录:
- 下载
kingbase8-8.6.0.jar(Maven 坐标cn.com.kingbase:kingbase8:8.6.0); - 将其复制到
$SEATUNNEL_HOME/plugins/jdbc/lib/目录下;
cp kingbase8-8.6.0.jar $SEATUNNEL_HOME/plugins/jdbc/lib/$SEATUNNEL_HOME即 SeaTunnel 的安装工作目录。如果不放置驱动,任务启动时会因找不到com.kingbase8.Driver而报ClassNotFoundException。
驱动类的自动识别原理
连接器通过 KingbaseDialectFactory 注册 Kingbase 方言,其核心判定逻辑是 URL 前缀匹配:
@Override public boolean acceptsURL(String url) { return url.startsWith("jdbc:kingbase8:"); }即只要 Sink 配置中的url以jdbc:kingbase8:开头,连接器就会自动加载 Kingbase 专属的 Dialect(dialectFactoryName()返回DatabaseIdentifier.KINGBASE,即"KingBase",定义见 DatabaseIdentifier.java)。这也是配置中url与driver必须正确填写的根本原因。
数据类型映射
Kingbase 数据类型与 SeaTunnel 数据类型的映射关系如下表(来自官方文档):
| Kingbase 数据类型 | SeaTunnel 数据类型 |
|---|---|
| BOOL | BOOLEAN |
| INT2 | SHORT |
| SMALLSERIAL / SERIAL / INT4 | INT |
| INT8 / BIGSERIAL | BIGINT |
| FLOAT4 | FLOAT |
| FLOAT8 | DOUBLE |
| NUMERIC | DECIMAL(获取指定列的指定列大小, 获取指定列小数点右边的位数) |
| BPCHAR / CHARACTER / VARCHAR / TEXT | STRING |
| TIMESTAMP | LOCALDATETIME |
| TIME | LOCALTIME |
| DATE | LOCALDATE |
| 其他数据类型 | 暂不支持 |
底层类型转换器的实现细节
上述映射在源码中由 KingbaseTypeConverter 实现。值得关注的是:
- 该类继承自
PostgresTypeConverter,因为 Kingbase 与 PostgreSQL 兼容,绝大多数类型(BOOL、INT2、FLOAT4、TIMESTAMP、TEXT 等)直接复用 PostgreSQL 的转换逻辑; - 在 PostgreSQL 无法处理时,会针对 Kingbase 的多模式兼容特性做额外处理(源码注释引用了 Kingbase 官方数据类型文档):例如 MySQL 兼容模式的
INT/MEDIUMINT/DATETIME/TINYBLOB等类型、Oracle 兼容模式的NUMBER/FLOAT/VARCHAR2/ROWID等类型,以及 Kingbase 特有的TINYINT(映射为 BYTE)、MONEY(映射为DECIMAL(38,18))、BLOB(映射为字节数组)、CLOB(映射为 STRING)、BIT(按BIT(M) -> BYTE(M/8)向上取整)等; - 类型读取侧的 KingbaseTypeMapper 则从
ResultSetMetaData提取列名、原生类型名、可空性、精度与小数位,统一交给KingbaseTypeConverter.INSTANCE完成 SeaTunnel 类型转换。
也就是说,官方文档表格是最小公共集合,实际运行时还会根据 Kingbase 的兼容模式(PG/MySQL/Oracle)自动扩展支持更多类型。
Sink 选项详解
以下为 Kingbase Sink 的全部配置参数(来自官方文档,并补充了源码 JdbcSinkOptions.java 中的默认值佐证):
| 参数名 | 类型 | 必须 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL,如jdbc:kingbase8://localhost:54321/db_test |
| driver | String | 是 | - | JDBC 驱动类名,Kingbase 固定为com.kingbase8.Driver |
| username | String | 否 | - | 连接用户名;旧配置名user仍可作为兼容写法使用 |
| password | String | 否 | - | 连接密码 |
| query | String | 否 | - | 使用该 SQL 将上游数据写入数据库(如INSERT ...),query优先级更高 |
| database | String | 否 | - | 配合table自动生成 SQL 并写入;与query互斥且优先级更高 |
| table | String | 否 | - | 配合database自动生成 SQL 并写入;与query互斥且优先级更高 |
| primary_keys | Array | 否 | - | 自动生成 SQL 时支持insert、delete、update等操作 |
| connection_check_timeout_sec | Int | 否 | 30 | 等待连接校验的数据库操作完成的时间(秒) |
| max_retries | Int | 否 | 0 | 提交失败(executeBatch)的重试次数 |
| batch_size | Int | 否 | 1000 | 批量写入缓冲,达到batch_size条或checkpoint.interval时间时刷新 |
| batch_interval_ms | Long | 否 | 0 | 定时刷新间隔(毫秒);0表示关闭,大于0时每次写入检查是否超过该间隔,超过则同步刷新 |
| is_exactly_once | Boolean | 否 | false | 是否启用精确一次(XA 事务),需同时设置xa_data_source_class_name;Kingbase 不支持 |
| generate_sink_sql | Boolean | 否 | false | 根据目标表自动生成 SQL |
| xa_data_source_class_name | String | 否 | - | XA 数据源类名;Kingbase 不支持 |
| max_commit_attempts | Int | 否 | 3 | 事务提交失败的重试次数 |
| transaction_timeout_sec | Int | 否 | -1 | 事务超时时间,默认 -1 永不超时;设置超时可能影响精确一次语义 |
| auto_commit | Boolean | 否 | true | 默认启用自动事务提交 |
| enable_upsert | Boolean | 否 | true | 存在 primary_keys 时启用 upsert;任务无重复数据时设false可加速导入 |
| common-options | - | 否 | - | Sink 插件通用参数,见 Sink 通用选项 |
选项要点与源码印证
enable_upsert的底层实现:当指定了primary_keys时,KingbaseDialect 的getUpsertStatement会生成 PostgreSQL 风格的INSERT ... ON CONFLICT (pk) DO UPDATE SET col=EXCLUDED.col语句,这是 Kingbase 与 PostgreSQL 兼容的重要体现;- 标识符引用:
quoteIdentifier使用双引号包裹标识符(如"table"、"column"),并支持schema.table形式的点分解析;同时可通过field_ide选项控制字段大小写转换策略; - 并发提示:官方文档特别提示——若未设置
partition_column,任务将以单并发运行;设置了partition_column则会根据任务并发度并行执行。
任务示例
运行以下任务前,需要先完成两项准备:
- 在 Kingbase 中创建目标数据库与表(如数据库
test、表test_table); - 若尚未安装部署 SeaTunnel,请先参考 安装 SeaTunnel 完成部署,再按照 使用 SeaTunnel 引擎快速开始 运行作业。
示例一:简单写入(自定义 SQL)
该示例通过 FakeSource 自动生成 16 行数据(row.num=16),每行 12 个字段,最终写入 Kingbase 的test_table表:
# 定义运行时环境 env { parallelism = 1 job.mode = "BATCH" } source { # 这是一个示例源插件,仅用于测试和演示源插件功能 FakeSource { parallelism = 1 plugin_output = "fake" row.num = 16 schema = { fields { c_string = string c_boolean = boolean c_tinyint = tinyint c_smallint = smallint c_int = int c_bigint = bigint c_float = float c_double = double c_decimal = "decimal(30, 8)" c_date = date c_time = time c_timestamp = timestamp } } } } transform { } sink { jdbc { url = "jdbc:kingbase8://127.0.0.1:54321/dbname" driver = "com.kingbase8.Driver" username = "root" password = "123456" query = "insert into test_table(c_string,c_boolean,c_tinyint,c_smallint,c_int,c_bigint,c_float,c_double,c_decimal,c_date,c_time,c_timestamp) values(?,?,?,?,?,?,?,?,?,?,?,?)" } }要点说明:
- 使用
query时,SQL 中的?占位符数量与顺序必须与上游字段一一对应; - FakeSource 的
c_tinyint字段在 Kingbase 侧对应TINYINT类型(通过KingbaseTypeConverter的 KB_TINYINT 分支映射为 SeaTunnel BYTE); - 若需了解源插件完整列表,可查阅 connector 文档。
示例二:自动生成 Sink SQL
无需手写复杂的 INSERT 语句,配置数据库名与表名即可自动生成写入 SQL:
sink { jdbc { url = "jdbc:kingbase8://127.0.0.1:54321/dbname" driver = "com.kingbase8.Driver" username = "root" password = "123456" # 根据数据库表名自动生成 sql 语句 generate_sink_sql = true database = test table = test_table } }示例三:写入带 Schema 的表
使用自定义query写入时,占位符数量必须与上游字段数量一致。Kingbase 带 schema 的表可写成public.table_name:
sink { Jdbc { driver = "com.kingbase8.Driver" url = "jdbc:kingbase8://localhost:54321/test" user = "SYSTEM" password = "123456" query = "INSERT INTO public.e2e_table_sink (c1, c2, c3) VALUES (?, ?, ?)" } }该示例同时演示了旧配置名user(等价于username)的兼容用法,以及schema.table形式的表名写法。在源码层面,KingbaseDialect.quoteIdentifier会按.拆分标识符并分别加双引号引用,因此public.e2e_table_sink会被正确解析为"public"."e2e_table_sink"。
高级能力:Kingbase Catalog 与自动建表
除 Sink 写入外,连接器还提供了完整的 Catalog 实现(源码位于 catalog/kingbase),可用于元数据管理与目标表自动创建:
- KingbaseCatalog:通过系统表
sys_class、sys_namespace、sys_attribute等查询列名、类型、长度、精度、默认值与注释,并自动排除INFORMATION_SCHEMA、SYSAUDIT、SYSLOGICAL、SYS_CATALOG、SYS_HM、XLOG_RECORD_READ等系统 schema; - KingbaseCreateTableSqlBuilder:负责生成建表 DDL,支持主键约束(超长主键名自动截断并追加随机后缀)、列注释、
WITH (fillfactor=...)与TABLESPACE "..."表选项; - KingbaseDialect.validateTableOptions:对
table_options做严格校验——fillfactor必须是 10~100 的整数,tablespace不允许包含引号、换行与分号等危险字符,否则直接抛出配置校验异常。
这些能力与schema_save_mode、data_save_mode(源码见 JdbcSinkOptions.java)配合,可实现"自动建表 + 增量追加/覆盖写"的整库同步场景。
变更日志
Kingbase 连接器的演进记录请查阅 connector-jdbc 变更日志。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考