跳到主要内容
版本:3.0.0

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-namesList否从 url 中解析出的数据库要监控的数据库名。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.modeEnum否INITIALDB2 CDC 的可选启动模式,有效值为 initial、earliest 和 latest。
stop.modeEnum否NEVERDB2 CDC 的可选停止模式,有效值为 never。
incremental.parallelismInteger否1增量阶段中并行读取器的数量。
snapshot.split.sizeInteger否8096表快照的分割大小。
snapshot.fetch.sizeInteger否1024读取表快照时每次轮询的最大获取大小。
server-time-zoneString否UTC数据库服务器中的会话时区。
connect.timeout.msDuration否30s连接器尝试连接到数据库服务器后,在超时之前等待的最长时间。
connect.max-retriesInteger否3连接器重试建立数据库连接的最大次数。
connection.pool.sizeInteger否20连接池大小。
chunk-key.even-distribution.factor.upper-boundDouble否100用于判断分块键是否均匀分布的上界。
chunk-key.even-distribution.factor.lower-boundDouble否0.05用于判断分块键是否均匀分布的下界。
sample-sharding.thresholdint否1000分块键分布不均时触发采样分片策略的估计分片数阈值。
inverse-sampling.rateint否1000采样分片策略使用的采样率倒数。
exactly_onceBoolean否false启用初始快照切换增量阶段时的精确一次语义。
debezium.*config否-透传给 Debezium DB2 连接器的配置项。
formatEnum否DEFAULT可选输出格式,有效值为 DEFAULT 和 COMPATIBLE_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
[Bug][Connector-V2][CDC] Prune removed tables from restored incremental splits (#11271)https://github.com/apache/seatunnel/commit/af0a647d23.0.0
[Feature][Connector-V2] Add DB2 CDC source connector (#10780)https://github.com/apache/seatunnel/commit/a9f69848a3.0.0