SeaTunnel DB2 Sink 实战:通过 Jdbc 连接器写入 DB2,掌握 XA 精确一次与 MERGE Upsert
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文以 SeaTunnel 官方 DB2 Sink 文档为主线,完整覆盖 JDBC DB2 Sink 连接器的依赖安装、全部配置项、批/流写入、XA 事务精确一次语义以及基于主键的 Upsert 写入四类实战场景,并结合开源仓库中connector-jdbc模块的 DB2 方言源码(方言工厂、类型转换器、MERGE 语句生成、Catalog 建表)与端到端测试用例,说明每个配置项在底层是如何生效的,帮助你在生产环境中把数据可靠、幂等地写入 DB2。
连接器定位与支持引擎
DB2 Sink 本质上是 SeaTunnel 的Jdbc 通用连接器在 DB2 数据库上的方言实现:配置中sink { Jdbc { ... } }的插件名是Jdbc,SeaTunnel 根据url的前缀(jdbc:db2:)自动路由到 DB2 方言。
该连接器支持以下引擎:
- Spark
- Flink
- SeaTunnel Zeta
从源码结构看,方言的路由机制由工厂类 DB2DialectFactory 完成,它通过@AutoService注册,并在acceptsURL方法中以url.startsWith("jdbc:db2:")判定是否接管当前连接:
@AutoService(JdbcDialectFactory.class) public class DB2DialectFactory implements JdbcDialectFactory { @Override public boolean acceptsURL(String url) { return url.startsWith("jdbc:db2:"); } // ... }这意味着你不需要(也不应该)在配置中显式指定数据库类型,只要 URL 写对,方言、行转换器、类型映射器都会自动装配。
关键特性
DB2 Sink 具备以下能力:
- 批处理(Batch)
- 流处理(Stream)
- 精确一次(Exactly-Once,基于 XA 事务)
- CDC 写入(通过主键 upsert / merge SQL 实现变更应用)
- 支持多表写入(配合
table占位符按上游表名路由) - 定时刷新(
batch_interval_ms定时 flush)
其中精确一次需要同时启用is_exactly_once = true并配置数据库对应的xa_data_source_class_name;CDC 语义则依赖generate_sink_sql = true+primary_keys生成的MERGEupsert 语句。
使用依赖(JDBC 驱动安装)
DB2 官方 JDBC 驱动(Maven 坐标com.ibm.db2.jcc:db2jcc)不随 SeaTunnel 发行包内置,需要手动放置:
| 运行引擎 | 驱动 JAR 放置目录 |
|---|---|
| Spark / Flink | ${SEATUNNEL_HOME}/plugins/ |
| SeaTunnel Zeta | ${SEATUNNEL_HOME}/lib/ |
驱动类名为com.ibm.db2.jcc.DB2Driver,XA 数据源类名为com.ibm.db2.jcc.DB2XADataSource,两者均来自db2jcc这个 JAR。
支持的数据源信息
| 数据库 | 驱动 | URL 格式 |
|---|---|---|
| DB2 | com.ibm.db2.jcc.DB2Driver | jdbc:db2://127.0.0.1:50000/dbname |
注意不同版本的db2jcc依赖对应的驱动类路径一致(均为com.ibm.db2.jcc.DB2Driver),可按项目环境选择合适的 JAR 版本;JDBC URL 已经绑定到具体数据库(catalog),因此后续 SQL 中表名是否需要带 schema 前缀由方言单独处理(下文"表标识符处理"一节说明)。
数据类型映射
DB2 与 SeaTunnel 之间的类型映射关系如下(源自 DB2 Sink 文档):
| DB2 数据类型 | SeaTunnel 数据类型 |
|---|---|
| BOOLEAN | BOOLEAN |
| SMALLINT | SHORT |
| INT / INTEGER | INTEGER |
| BIGINT | LONG |
| DECIMAL / DEC / NUMERIC / NUM | DECIMAL(38,18) |
| REAL | FLOAT |
| FLOAT / DOUBLE / DOUBLE PRECISION / DECFLOAT | DOUBLE |
| CHAR / VARCHAR / LONG VARCHAR / CLOB / GRAPHIC / VARGRAPHIC / LONG VARGRAPHIC / DBCLOB | STRING |
| BLOB | BYTES |
| DATE | DATE |
| TIME | TIME |
| TIMESTAMP | TIMESTAMP |
| ROWID / XML | 暂不支持 |
源码中的映射实现细节
上述映射的运行时实现位于 DB2TypeConverter,除了上面的基本对应关系,源码中还定义了若干对建表(Catalog 自动建表)至关重要的边界约束,建议在写 DDL 时提前了解:
- DECIMAL 精度上限 31(
MAX_PRECISION = 31,MAX_SCALE = 30)。当 SeaTunnel 侧 DECIMAL 精度超过 31 时,转换器会截断精度、相应收缩 scale 并打 warn 日志,而不是直接报错。 - TIMESTAMP 精度上限 12(
MAX_TIMESTAMP_SCALE = 12),即 DB2 最大支持到皮秒级;超出会被截断为TIMESTAMP(12)。 - STRING 长度分级:长度 ≤ 255 建
CHAR(n),≤ 32672 建VARCHAR(n),更大建CLOB;未指定长度时默认VARCHAR(32672)。 - BYTES 长度分级:未指定长度默认
VARBINARY(32672);≤ 255 建BINARY(n),≤ 32672 建VARBINARY(n),更大建BLOB。 - 文档标注 ROWID、XML 在 Sink 路径"暂不支持",但从
DB2TypeConverter.convert的 switch 分支看,XML 在读取侧可被映射为 STRING;若遇到 ROWID 列会抛出类型转换错误。实际使用时建议避开 ROWID 列。
类型识别入口 DB2TypeMapper 会从ResultSetMetaData中读取列名、原生类型名、是否可空、precision/scale 等元数据,再委托DB2TypeConverter完成到 SeaTunnelColumn的转换——这也是generate_sink_sql/ Catalog 能"按源表结构生成目标结构"的元数据来源。
配置项详解
以下选项表完整继承自官方 DB2 文档,默认值均与 JdbcSinkOptions 源码中的定义一致:
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL,例如jdbc:db2://127.0.0.1:50000/dbname |
| driver | String | 是 | - | JDBC 驱动类名,DB2 使用com.ibm.db2.jcc.DB2Driver |
| username | String | 否 | - | DB2 用户名 |
| password | String | 否 | - | DB2 密码 |
| query | String | 否 | - | 写入上游数据的 SQL。优先级高于database/table自动生成的 SQL;设置后会关闭基于目录(Catalog)的优化(无法生成MERGEupsert) |
| database | String | 否 | - | 数据库名。generate_sink_sql = true时与table一起用于生成INSERT/MERGESQL;与query互斥,同时设置时query优先 |
| table | String | 否 | - | 目标表名。与database一起配合generate_sink_sql生成写入语句 |
| primary_keys | Array | 否 | - | 主键列。generate_sink_sql = true且enable_upsert = true时用于构建MERGEupsert 语句 |
| connection_check_timeout_sec | Int | 否 | 30 | 连接校验超时时间(秒) |
| max_retries | Int | 否 | 0 | executeBatch失败的重试次数 |
| batch_size | Int | 否 | 1000 | 触发 flush 的缓冲行数;同时在checkpoint.interval时也会 flush |
| batch_interval_ms | Long | 否 | 0 | 两次 flush 之间的最大时间间隔(毫秒)。0表示关闭按时间间隔的 flush |
| is_exactly_once | Boolean | 否 | false | 是否启用基于 XA 的精确一次;启用时必须设置xa_data_source_class_name |
| generate_sink_sql | Boolean | 否 | false | 基于database/table/primary_keys自动生成INSERT或MERGESQL,而不是手动提供query |
| xa_data_source_class_name | String | 否 | - | XA 数据源类名,DB2 使用com.ibm.db2.jcc.DB2XADataSource |
| max_commit_attempts | Int | 否 | 3 | 事务提交失败的重试次数 |
| transaction_timeout_sec | Int | 否 | -1 | 事务超时时间(秒),-1表示永不超时;设置超时可能会影响精确一次 |
| auto_commit | Boolean | 否 | true | 是否自动提交每个批次 |
| properties | Map | 否 | - | 额外的 JDBC 连接参数。properties与url包含相同键时优先级由驱动决定 |
| common-options | - | 否 | - | Sink 插件通用参数,详见 Sink 通用选项 |
| enable_upsert | Boolean | 否 | true | 在primary_keys已配置且generate_sink_sql = true时,生成MERGEupsert 语句;若输入无重复主键,可设为false使用更快的纯插入路径 |
几点源码层面的补充说明:
batch_interval_ms在 JdbcSinkOptions 中的注释明确了其触发方式:0为关闭(默认);大于 0 时每次写入记录都会检查距上次 flush 的时间,超过间隔即同步 flush。流式低吞吐场景下,这个参数可以避免数据长时间滞留缓冲。table支持${table_name}占位符实现多表写入(源码中已废弃独立的tablePrefix/tableSuffix选项,统一改用table = "prefix_${table_name}_suffix"的形式)。generate_sink_sql默认false,即默认必须显式提供query;开启它才能走自动生成 SQL 与 Catalog 优化路径。
任务示例
示例一:手动 SQL 简单写入
从FakeSource读取 16 行数据插入 DB2 的test_table:
env { parallelism = 1 job.mode = "BATCH" } source { FakeSource { parallelism = 1 plugin_output = "fake" row.num = 16 schema = { fields { name = "string" age = "int" } } } } sink { Jdbc { url = "jdbc:db2://127.0.0.1:50000/dbname" driver = "com.ibm.db2.jcc.DB2Driver" username = "db2inst1" password = "123456" query = "insert into test_table(name, age) values(?, ?)" } }运行作业前,请先在 DB2 中创建目标数据库和表。这是最"朴素"的用法:驱动执行PreparedStatement批量executeBatch,batch_size(默认 1000 行)满时提交一次。
示例二:自动生成 Sink SQL
不写INSERT语句,让 SeaTunnel 根据database和table自动生成:
sink { Jdbc { url = "jdbc:db2://127.0.0.1:50000/dbname" driver = "com.ibm.db2.jcc.DB2Driver" username = "db2inst1" password = "123456" generate_sink_sql = true database = test table = test_table } }开启generate_sink_sql后,连接器会通过 Catalog 机制获取目标表列结构(或按上游 Schema 自动建表),生成带列名的INSERT INTO ... VALUES (?, ?, ...),避免手写 SQL 与表结构漂移。注意:此时不能再同时指望query生效——query优先级虽更高,但一旦设置query,基于 Catalog 的优化(含自动建表、upsert)会整体关闭。
示例三:基于 XA 事务的精确一次
sink { Jdbc { url = "jdbc:db2://127.0.0.1:50000/dbname" driver = "com.ibm.db2.jcc.DB2Driver" username = "db2inst1" password = "123456" query = "insert into test_table(name, age) values(?, ?)" max_retries = 0 is_exactly_once = true xa_data_source_class_name = "com.ibm.db2.jcc.DB2XADataSource" } }启用is_exactly_once = true后,写入路径切换到 XA 两阶段提交:只有在 checkpoint 成功且两阶段提交都成功时数据才对下游可见;作业失败回滚后重跑不会留下"半提交"的重复或残缺批次。配套的容错参数:
max_retries:批次执行失败的重试次数(精确一次模式下通常保持0,失败交由事务回滚与作业重试处理);max_commit_attempts(默认 3):阶段二"提交"失败时的重试次数;transaction_timeout_sec(默认-1不超时):设置过短的事务超时可能在长 checkpoint 间隔下引发回滚,反而破坏精确一次,流式作业中要谨慎调整。
对应实现见 JdbcExactlyOnceSinkWriter。需要说明的前提:XA 精确一次依赖目标数据库支持两阶段提交并配置好 XA 数据源,DB2 上即DB2XADataSource;同时 XA 连接通常不允许 autocommit,启用该模式时auto_commit保持默认即可,由事务管理器接管提交时机。
示例四:自动生成 SQL 的 Upsert(CDC 场景)
当generate_sink_sql = true且设置了primary_keys时,DB2 通过生成的MERGE语句完成 upsert 写入;如果上游数据只包含新增(无更新/无重复主键),可设置enable_upsert = false走更快的纯插入路径:
sink { Jdbc { url = "jdbc:db2://127.0.0.1:50000/E2E" driver = "com.ibm.db2.jcc.DB2Driver" username = "db2inst1" password = "123456" database = "E2E" table = "SINK" generate_sink_sql = true enable_upsert = true primary_keys = ["C_INT"] } }这与仓库端到端测试 jdbc_db2_source_and_sink_upsert.conf 的配置完全同构:从E2E.SOURCE表读出数据,按主键C_INT合并写入E2E.SINK表,由 JdbcDb2UpsertIT 断言写入结果(纯插入场景的基线用例见 JdbcDb2IT)。这也是 CDC 管道(如 DB2 CDC Source 或其他变更源)落库时消除主键冲突、保证幂等的推荐写法。
源码深度解析:DB2 的 MERGE 语句是怎么生成的
DB2 没有 MySQL 那样的ON DUPLICATE KEY UPDATE或 PostgreSQL 的ON CONFLICT子句,SeaTunnel 采用 DB2 标准的MERGE INTO ... USING (VALUES (...))语法实现 upsert。生成逻辑集中在 DB2Dialect.getUpsertStatement,按主键与列名拼装出形如:
MERGE INTO "DB"."TABLE" AS target USING (VALUES (?, ?)) AS source ("COL1", "COL2") ON target."COL1" = source."COL1" AND target."COL2" = source."COL2" WHEN MATCHED AND (target."COL1" <> source."COL1" OR target."COL2" <> source."COL2") THEN UPDATE SET target."COL1" = source."COL1", target."COL2" = source."COL2" WHEN NOT MATCHED THEN INSERT ("COL1", "COL2") VALUES (source."COL1", source."COL2")其构造分为五步,与源码逐一对应:
- USING 子句:把每行数据包装成
VALUES (?, ?, ...)派生表source,列名列表由字段名经quoteIdentifier加双引号生成; - ON 子句:所有主键列做
target.pk = source.pk的 AND 连接,即匹配条件; - WHEN MATCHED 守卫:追加
target.col <> source.col OR ...条件——只有源数据与目标行确有差异时才执行 UPDATE,可以避免无意义更新; - UPDATE SET:把非匹配差异列全部从
source回写到target; - WHEN NOT MATCHED:主键不存在时执行
INSERT,完成"插入"半边。
另外两个与 DB2 特性相关的实现细节值得关注:
表标识符处理。DB2Dialect.tableIdentifier中有一段针对 DB2 的专门处理:JDBC URL 已经把连接绑定到某个数据库(catalog),如果表名本身携带 schema(如SCHEMA.TABLE),再前置database会拼出非法的三段式标识符,因此带点的表名只会被整体加引号,不带点时才拼成"db"."table":
// DB2 connections are already bound to a database by the JDBC URL. When the table name // carries a schema, prefixing the configured database would generate an invalid // catalog.schema.table identifier. if (tableName.contains(".")) { return quoteIdentifier(tableName); } return quoteIdentifier(database) + "." + quoteIdentifier(tableName);这解释了为什么 E2E 配置里table = "SINK"而查询 SQL 中写的是"E2E".SOURCE——URL 中的E2E与database选项共同决定了最终 SQL 里的表定位方式。
方言标识与引号风格。dualTable()返回FROM SYSIBM.SYSDUMMY1,即 DB2 执行SELECT 1这类探针查询时使用的"哑元表",用于连接可用性检查等场景;标识符统一用双引号包裹(quoteIdentifier),并对复合标识符a.b.c逐段加引号。
Catalog:自动建表与元数据能力
除了方言,DB2 还支持 SeaTunnel 的 Catalog 抽象:DB2CatalogFactory 以DB2作为 factory identifier,从配置中读取url/username/password/schema/driver构造 DB2Catalog。在开启generate_sink_sql且目标表不存在时,Catalog 会依据上游 Schema 调用DB2TypeConverter.reconvert(上文提到的"SeaTunnel 类型 → DB2 DDL 类型"反向映射,含 DECIMAL 截断、CHAR/VARCHAR/CLOB 长度分级等规则)生成建表语句,从而实现"源表什么样,目标表就建成什么样"。
生产使用建议
结合文档与源码,可以给出几条可直接落地的建议:
- 纯追加日志类数据:用示例一/二的
INSERT路径 + 默认auto_commit,必要时配batch_interval_ms控制延迟; - 流式 + 故障重跑必须无重复:启用示例三的 XA 精确一次,
max_commit_attempts保持默认 3,谨慎设置transaction_timeout_sec; - CDC / 含更新语义的数据:用示例四的
generate_sink_sql + primary_keys + enable_upsert = true,让 MERGE 语句兜住主键冲突;若确认流里只有INSERT事件(如 Debezium 过滤后),设enable_upsert = false可省掉 MERGE 的匹配开销; - 多表写入:利用
table = "prefix_${table_name}"占位符,一份配置承接多个上游表; - 类型核对:目标库为 DB2 时,留意 DECIMAL 精度(>31 截断)、TIMESTAMP 精度(>12 截断)的自动降级行为,避免精度损失超出业务预期。
变更记录
Jdbc 连接器的历史版本变更记录,请参见 Jdbc 连接器变更日志;DB2 方言相关的实现代码可继续深入 connector-jdbc 模块 的internal/dialect/db2与catalog/db2两个包查看。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考