SeaTunnel Kingbase Sink 连接器完全指南:JDBC 配置、类型映射与实战写入

发布时间:2026/9/18 21:09:59
SeaTunnel Kingbase Sink 连接器完全指南:JDBC 配置、类型映射与实战写入
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支持引擎SparkFlinkSeaTunnel Zeta关键特性特性支持情况精确一次Exactly Once❌ 不支持CDC❌ 不支持定时刷新✅ 支持需要特别说明的是官方文档明确指出连接器依赖XA 事务来保证精确一次语义因此只有支持XA 事务的数据库才能启用精确一次通过is_exactly_oncetrue而Kingbase 目前不支持 XA 事务所以该特性在 Kingbase 上不可用。与此对应的xa_data_source_class_name选项同样不适用于 Kingbase。支持的数据源信息数据源支持的版本驱动URLMavenKingbase8.6com.kingbase8.Driverjdbc:kingbase8://localhost:54321/db_testkingbase8-8.6.0.jar从仓库的 connector-jdbc/pom.xml 可以看到JDBC 连接器模块已将cn.com.kingbase:kingbase8以8.6.0版本声明为依赖kingbase8.version8.6.0/kingbase8.version与实际支持的版本一致。数据库依赖安装使用该连接器前需要下载对应的 JDBC 驱动 jar 并放入 SeaTunnel 的插件目录下载kingbase8-8.6.0.jarMaven 坐标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 专属的 DialectdialectFactoryName()返回DatabaseIdentifier.KINGBASE即KingBase定义见 DatabaseIdentifier.java。这也是配置中url与driver必须正确填写的根本原因。数据类型映射Kingbase 数据类型与 SeaTunnel 数据类型的映射关系如下表来自官方文档Kingbase 数据类型SeaTunnel 数据类型BOOLBOOLEANINT2SHORTSMALLSERIAL / SERIAL / INT4INTINT8 / BIGSERIALBIGINTFLOAT4FLOATFLOAT8DOUBLENUMERICDECIMAL(获取指定列的指定列大小, 获取指定列小数点右边的位数)BPCHAR / CHARACTER / VARCHAR / TEXTSTRINGTIMESTAMPLOCALDATETIMETIMELOCALTIMEDATELOCALDATE其他数据类型暂不支持底层类型转换器的实现细节上述映射在源码中由 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 中的默认值佐证参数名类型必须默认值描述urlString是-JDBC 连接 URL如jdbc:kingbase8://localhost:54321/db_testdriverString是-JDBC 驱动类名Kingbase 固定为com.kingbase8.DriverusernameString否-连接用户名旧配置名user仍可作为兼容写法使用passwordString否-连接密码queryString否-使用该 SQL 将上游数据写入数据库如INSERT ...query优先级更高databaseString否-配合table自动生成 SQL 并写入与query互斥且优先级更高tableString否-配合database自动生成 SQL 并写入与query互斥且优先级更高primary_keysArray否-自动生成 SQL 时支持insert、delete、update等操作connection_check_timeout_secInt否30等待连接校验的数据库操作完成的时间秒max_retriesInt否0提交失败executeBatch的重试次数batch_sizeInt否1000批量写入缓冲达到batch_size条或checkpoint.interval时间时刷新batch_interval_msLong否0定时刷新间隔毫秒0表示关闭大于0时每次写入检查是否超过该间隔超过则同步刷新is_exactly_onceBoolean否false是否启用精确一次XA 事务需同时设置xa_data_source_class_nameKingbase 不支持generate_sink_sqlBoolean否false根据目标表自动生成 SQLxa_data_source_class_nameString否-XA 数据源类名Kingbase 不支持max_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务超时时间默认 -1 永不超时设置超时可能影响精确一次语义auto_commitBoolean否true默认启用自动事务提交enable_upsertBoolean否true存在 primary_keys 时启用 upsert任务无重复数据时设false可加速导入common-options-否-Sink 插件通用参数见 Sink 通用选项选项要点与源码印证enable_upsert的底层实现当指定了primary_keys时KingbaseDialect 的getUpsertStatement会生成 PostgreSQL 风格的INSERT ... ON CONFLICT (pk) DO UPDATE SET colEXCLUDED.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.num16每行 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 语句配置数据库名与表名即可自动生成写入 SQLsink { 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_namesink { 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等系统 schemaKingbaseCreateTableSqlBuilder负责生成建表 DDL支持主键约束超长主键名自动截断并追加随机后缀、列注释、WITH (fillfactor...)与TABLESPACE ...表选项KingbaseDialect.validateTableOptions对table_options做严格校验——fillfactor必须是 10100 的整数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),仅供参考