跳到主要内容
版本:Next

DB2 CDC

DB2 CDC 源连接器

支持 DB2 版本

  • DB2 LUW 11.5 或 Debezium DB2 连接器支持的更高版本

支持的引擎

SeaTunnel Zeta
Flink

主要功能

描述

DB2 CDC 连接器可以读取已启用 capture mode 的 DB2 表的快照数据和增量数据。连接器内部使用 Debezium DB2,在初始快照完成后会继续读取已提交的 INSERT、UPDATE 和 DELETE 变更。

支持的数据源信息

数据源支持版本驱动UrlMaven
DB2DB2 LUW 11.5 或 Debezium DB2 连接器支持的更高版本com.ibm.db2.jcc.DB2Driverjdbc:db2://127.0.0.1:50000/testdbhttps://mvnrepository.com/artifact/com.ibm.db2.jcc/db2jcc

使用依赖

安装 Jdbc 驱动

  1. 你需要确保 DB2 JDBC 驱动 jar 包 已经放置在 ${SEATUNNEL_HOME}/plugins/ 目录中。

对于 SeaTunnel Zeta 引擎

  1. 你需要确保 DB2 JDBC 驱动 jar 包 已经放置在 ${SEATUNNEL_HOME}/lib/ 目录中。

数据类型映射

DB2 数据类型SeaTunnel 数据类型
BOOLEANBOOLEAN
SMALLINTSHORT
INT
INTEGER
INT
BIGINTBIGINT
DECIMAL
DEC
NUMERIC
NUM
DECIMAL
REALFLOAT
DOUBLE
DECFLOAT
DOUBLE
CHAR
CHARACTER
VARCHAR
LONG VARCHAR
CLOB
GRAPHIC
VARGRAPHIC
DBCLOB
XML
STRING
BINARY
VARBINARY
BLOB
BYTES
DATEDATE
TIMETIME
TIMESTAMPTIMESTAMP

数据源参数

名称类型是否必填默认值描述
usernameString-连接 DB2 时使用的用户名。
passwordString-连接 DB2 时使用的密码。
urlString-DB2 JDBC URL。URL 必须包含数据库名,例如 jdbc:db2://127.0.0.1:50000/testdb
database-namesListurl 中解析出的数据库要监控的数据库名。DB2 CDC 一个 source 监控一个数据库。
table-namesList未设置 table-pattern 时必填-要监控的表名,格式为 databaseName.schemaName.tableName,例如 testdb.DB2INST1.CUSTOMERS
table-patternString未设置 table-names 时必填-用于发现已开启 capture mode 的表的正则表达式。
table-names-configList-表配置列表。例如:[{"table": "testdb.DB2INST1.CUSTOMERS","primaryKeys": ["ID"],"snapshotSplitColumn": "ID"}]
startup.modeEnumINITIALDB2 CDC 的可选启动模式,有效值为 initialearliestlatest
stop.modeEnumNEVERDB2 CDC 的可选停止模式,有效值为 never
incremental.parallelismInteger1增量阶段中并行读取器的数量。
snapshot.split.sizeInteger8096表快照的分割大小。
snapshot.fetch.sizeInteger1024读取表快照时每次轮询的最大获取大小。
server-time-zoneStringUTC数据库服务器中的会话时区。
connect.timeout.msDuration30s连接器尝试连接到数据库服务器后,在超时之前等待的最长时间。
connect.max-retriesInteger3连接器重试建立数据库连接的最大次数。
connection.pool.sizeInteger20连接池大小。
chunk-key.even-distribution.factor.upper-boundDouble100用于判断分块键是否均匀分布的上界。
chunk-key.even-distribution.factor.lower-boundDouble0.05用于判断分块键是否均匀分布的下界。
sample-sharding.thresholdint1000分块键分布不均时触发采样分片策略的估计分片数阈值。
inverse-sampling.rateint1000采样分片策略使用的采样率倒数。
exactly_onceBooleanfalse启用初始快照切换增量阶段时的精确一次语义。
debezium.*config-透传给 Debezium DB2 连接器的配置项。
formatEnumDEFAULT可选输出格式,有效值为 DEFAULTCOMPATIBLE_DEBEZIUM_JSON
common-options-源插件通用参数,请参考 Source Common Options 获取详细信息。

启用 DB2 CDC

DB2 CDC 依赖 DB2 SQL replication 和 ASN capture tables。启用 capture 前,请先确认当前环境具备所需的 IBM replication 授权。运行 SeaTunnel 前,数据库管理员必须先把要读取的表加入 capture mode。可以使用 DB2 控制命令,也可以使用 Debezium 提供的管理 UDF。下面是常见的 UDF 流程:

VALUES ASNCDC.ASNCDCSERVICES('status','asncdc');
VALUES ASNCDC.ASNCDCSERVICES('start','asncdc');
CALL ASNCDC.ADDTABLE('DB2INST1', 'CUSTOMERS');
VALUES ASNCDC.ASNCDCSERVICES('reinit','asncdc');

完整的 DB2 服务端配置、权限和 ASN capture agent 配置请参考 Debezium DB2 连接器设置文档

任务示例

初始读取简单示例

该示例先读取初始快照,随后继续读取增量变更。

env {
parallelism = 1
job.mode = "STREAMING"
checkpoint.interval = 5000
}

source {
DB2-CDC {
plugin_output = "customers"
username = "db2inst1"
password = "db2inst1"
startup.mode = "initial"
database-names = ["testdb"]
table-names = ["testdb.DB2INST1.CUSTOMERS"]
url = "jdbc:db2://127.0.0.1:50000/testdb"
}
}

sink {
console {
plugin_input = "customers"
}
}

增量读取简单示例

该示例从最新 DB2 LSN 开始读取,并打印新产生的变更数据。

env {
parallelism = 1
job.mode = "STREAMING"
checkpoint.interval = 5000
}

source {
DB2-CDC {
plugin_output = "customers"
username = "db2inst1"
password = "db2inst1"
startup.mode = "latest"
database-names = ["testdb"]
table-names = ["testdb.DB2INST1.CUSTOMERS"]
url = "jdbc:db2://127.0.0.1:50000/testdb"
}
}

sink {
console {
plugin_input = "customers"
}
}

支持表的自定义主键

source {
DB2-CDC {
plugin_output = "customers"
username = "db2inst1"
password = "db2inst1"
startup.mode = "initial"
database-names = ["testdb"]
table-names = ["testdb.DB2INST1.CUSTOMERS"]
table-names-config = [
{
table = "testdb.DB2INST1.CUSTOMERS"
primaryKeys = ["ID"]
snapshotSplitColumn = "ID"
}
]
url = "jdbc:db2://127.0.0.1:50000/testdb"
}
}

变更日志

Change Log
ChangeCommitVersion