
Flink CDC Db2 Connector 实战指南从快照到增量日志的流式数据接入【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc导读Db2 CDC Connector 是 Flink CDCflink-cdc官方提供的流式数据接入组件用于从 IBM Db2 数据库同时读取存量快照数据snapshot与增量变更数据redo log。本文围绕 Db2 CDC 官方文档 展开完整覆盖依赖引入、Db2 服务端准备、Flink SQL 建表、全部连接器参数、启动模式、DataStream API 两种用法、无主键表支持、监控指标与数据类型映射并结合仓库源码说明其底层实现原理。读完本文你将能够在 Flink SQL 与 DataStream 两条技术路线上独立完成 Db2 数据的 CDC 实时同步。支持范围与版本前提Db2 CDC Connector 当前支持的数据库与驱动版本如下连接器支持的数据库驱动Db2-cdcDb211.5Db2 Driver11.5.0.0适用前提需要 Db2 开启归档日志并配置 SQL Replication即通过 Debezium 的 Db2 连接器机制订阅事务日志。连接器内部的实现依托 Debezium 嵌入式引擎仓库中 Db2SourceConfigFactory.java 直接以io.debezium.connector.db2.Db2Connector作为connector.class并注入 Flink 自定义的数据库历史记录实现EmbeddedFlinkDatabaseHistory来捕获 DDL 变更。依赖引入Maven 与 SQL Client JARMaven 依赖在 Maven/SBT 项目中引入连接器{{ artifact flink-connector-db2-cdc }}连接器模块的实际坐标可在仓库 flink-connector-db2-cdc/pom.xml 中确认其内部依赖flink-cdc-base、flink-connector-debezium与debezium-connector-db2。SQL Client JAR下载flink-sql-connector-db2-cdc包并放入FLINK_HOME/lib/目录即可在 Flink SQL Client 中使用。更多发布版本可在 Maven 中央仓库查询。⚠️ 重要IPLA 协议限制与 JDBC 驱动手动配置由于 Db2 连接器采用的 IPLA 协议与 Flink CDC 项目的许可证不兼容预构建的连接器 JAR 包中不包含 Db2 连接器需要手动补充以下依赖依赖名称说明com.ibm.db2.jcc:db2jcc:db2jcc4用于连接 Db2 数据库的 JDBC 驱动源码侧同样印证了这一点JDBC 驱动类被硬编码为com.ibm.db2.jcc.DB2Driver见 Db2SourceConfigFactory.java该驱动需由用户自行提供。准备 Db2 服务端连接器基于 Debezium Db2 Connector 实现Db2 服务端需要按照 Debezium Db2 Connector 的设置步骤 完成前置配置核心包括开启 Db2 归档日志模式为捕获增量变更提供 redo log 来源配置用于 CDC 的数据库用户并授予相应权限连接器在内部通过database.user/database.password建立会话。已知限制SQL Replication 不支持 BOOLEAN 类型Db2 的 SQL Replication 目前不支持BOOLEAN类型因此 Debezium 无法对包含BOOLEAN列的表执行 CDC 增量捕获只能对该类表做一次性快照读取。建议在业务设计中使用其他类型替代BOOLEAN。在 Flink SQL 中创建 Db2 CDC 表完整建表示例-- 每 3 秒触发一次 checkpoint Flink SQL SET execution.checkpointing.interval 3s; -- 在 Flink SQL 中注册 Db2 表 products Flink SQL CREATE TABLE products ( ID INT NOT NULL, NAME STRING, DESCRIPTION STRING, WEIGHT DECIMAL(10,3), PRIMARY KEY(ID) NOT ENFORCED ) WITH ( connector db2-cdc, hostname localhost, port 50000, username root, password 123456, database-name mydb, table-name myschema.products); -- 读取 products 表的快照数据与 redo log 增量数据 Flink SQL SELECT * FROM products;关键注意点主键声明PRIMARY KEY(ID) NOT ENFORCED中的NOT ENFORCED表示 Flink 不校验主键唯一性仅用于语义声明更新/删除操作的去重与下游幂等依赖它。在无主键表场景下必须显式配置scan.incremental.snapshot.chunk.key-column详见下文。table-name格式myschema.products即schema名.表名形式schema 与表名之间以点分隔。checkpoint 是必须的CDC 增量语义依赖 checkpoint 保存读取位点建议生产环境按业务可接受的延迟设置合理的 checkpoint 间隔。Connector Options 全参数详解以下为 Db2 CDC 连接器支持的完整参数表与 Db2TableSourceFactory.java 中定义的可选/必选参数一一对应Option是否必填默认值类型说明connectorrequired(none)String指定使用的连接器此处固定为db2-cdchostnamerequired(none)StringDb2 数据库服务器的 IP 地址或主机名usernamerequired(none)String连接 Db2 数据库服务器使用的用户名passwordrequired(none)String连接 Db2 数据库服务器使用的密码database-namerequired(none)String要监控的 Db2 服务器数据库名table-namerequired(none)String要监控的 Db2 表名格式为schema.table例如db1.table1portoptional50000IntegerDb2 数据库服务器的端口号scan.startup.modeoptionalinitialStringDb2 CDC 消费者的启动模式合法值为initial与latest-offset详见「启动读取位置」server-time-zoneoptional(none)String数据库服务器的会话时区如Asia/Shanghai控制 Db2 的 TIMESTAMP 类型如何转换为 STRING未设置时使用ZoneId.systemDefault()判断服务端时区scan.incremental.snapshot.enabledoptionaltrueBoolean是否启用并行快照读取chunk-meta.group.sizeoptional1000Integerchunk 元数据的分组大小元数据规模超过该值时会被拆分为多个组chunk-key.even-distribution.factor.lower-boundoptional0.05dDoublechunk key 分布因子的下界分布因子用于判断表数据是否均匀分布均匀时采用均匀计算优化切分不均匀时触发查询式切分。分布因子 (MAX(id) - MIN(id) 1) / rowCountchunk-key.even-distribution.factor.upper-boundoptional1000.0dDoublechunk key 分布因子的上界计算方式同上scan.incremental.snapshot.chunk.key-columnoptional(none)String快照读取时用于拆分 chunk 的键列默认取主键第一列。可以使用非主键列但可能降低查询性能debezium.*optional(none)String透传给 Debezium 嵌入式引擎的属性例如debezium.snapshot.mode neverscan.incremental.close-idle-reader.enabledoptionalfalseBoolean是否在快照阶段结束时关闭空闲 reader。当设置execution.checkpointing.checkpoints-after-tasks-finish.enabled为 true 时要求 Flink 版本 ≥ 1.14Flink ≥ 1.15 该配置默认即为 true无需显式设置scan.incremental.snapshot.unbounded-chunk-first.enabledoptionaltrueBoolean快照读取阶段是否优先分配无界 chunk有助于降低对最大无界 chunk 做快照时 TaskManager 发生 OOM 的风险源码级补充文档未列出的基类参数除上述表格外工厂还支持以下来自flink-cdc-base的增量快照调优参数见 SourceOptions.java 与 JdbcSourceOptions.javascan.incremental.snapshot.chunk.size快照分片行数默认8096scan.snapshot.fetch.size快照阶段每次拉取的行数默认1024connect.timeout连接 Db2 的超时时间默认30 秒connect.max-retries获取连接的最大重试次数默认3connection.pool.size连接池大小默认20scan.incremental.snapshot.backfill.skip是否跳过快照阶段的 backfill默认false。工厂在启用并行快照enableParallelRead true时会对上述整数参数做合法性校验chunk size / fetch size / chunk-meta group size / connection pool size 必须大于 1connect max retries 必须大于 0分布因子上下界也必须落在合法区间内详见 Db2TableSourceFactory.java。关于scan.startup.mode的冲突警告scan.startup.mode的机制依赖 Debezium 的snapshot.mode配置二者不要同时使用。如果在建表 DDL 中同时指定了scan.startup.mode与debezium.snapshot.mode可能导致scan.startup.mode不生效。源码印证Db2SourceConfigFactory内部正是把INITIAL映射为snapshot.modeinitial、把LATEST_OFFSET映射为snapshot.modeschema_only见 Db2SourceConfigFactory.java。可用 Metadata虚拟列在表定义中可以声明以下只读VIRTUAL元数据列KeyDataType说明table_nameSTRING NOT NULL包含该行数据的表名schema_nameSTRING NOT NULL包含该行数据的 schema 名database_nameSTRING NOT NULL包含该行数据的数据库名op_tsTIMESTAMP_LTZ(3) NOT NULL数据库变更发生的时间若记录来自表快照而非变更流该值恒为 0使用示例在 CREATE TABLE 中追加CREATE TABLE products ( ID INT NOT NULL, NAME STRING, PRIMARY KEY(ID) NOT ENFORCED, db_name STRING METADATA FROM database_name VIRTUAL, tb_name STRING METADATA FROM table_name VIRTUAL, ts TIMESTAMP_LTZ(3) METADATA FROM op_ts VIRTUAL ) WITH ( ... );从源码看这四类元数据由 Db2ReadableMetaData.java 定义其取值均从 DebeziumSourceRecord的source结构体中解析如AbstractSourceInfo.TABLE_NAME_KEY、TIMESTAMP_KEY等其中op_ts以TimestampData.fromEpochMillis(...)转换为时间戳。核心功能特性启动读取位置Startup Reading Positionscan.startup.mode决定 Db2 CDC 消费者的启动行为initial默认首次启动时对受监控的数据库表执行一次全量快照随后继续读取最新的 redo log 增量数据latest-offset首次启动时不对受监控的表执行快照直接从 redo log 末尾开始读取即只消费连接器启动之后的变更。注意如上文所述该参数与debezium.snapshot.mode二选一使用。快照分片Chunk机制当scan.incremental.snapshot.enabledtrue默认时快照读取采用增量快照算法按 chunk key 将表切分为多个 chunk分发给多个并行 subtask 读取从而支持并行快照。切分决策依赖两个因素chunk key默认取主键第一列也可通过scan.incremental.snapshot.chunk.key-column指定数据均匀性通过分布因子(MAX(id) - MIN(id) 1) / rowCount判断数据是否均匀均匀时采用均匀计算优化不均匀时触发查询式切分对应chunk-key.even-distribution.factor.lower-bound/upper-bound两个参数。快照完成后再无缝切换到 redo log 增量读取期间产生的变更通过 backfill 机制合并回快照scan.incremental.snapshot.backfill.skip可跳过该过程但会牺牲一致性保障。无主键表支持自 3.4.0 起自版本 3.4.0 开始Db2 CDC 支持无主键表。使用无主键表必须通过scan.incremental.snapshot.chunk.key-column指定一个非空字段作为 chunk key。需要注意两点优先选择索引内的列若表存在索引尽量将scan.incremental.snapshot.chunk.key-column指定为索引包含的列可显著提升 select 查询速度处理语义取决于该列的行为若该列从不发生更新操作则保证exactly-once精确一次语义若该列会发生更新操作则只能保证at-least-once至少一次语义。此时建议在下游为主键并执行幂等操作upsert来保证数据正确性。⚠️ 警告主键表使用非主键列作为 chunk key 可能导致数据不一致对有主键的 Db2 表若把非主键列配置为scan.incremental.snapshot.chunk.key-column可能出现数据不一致。以下为官方文档给出的问题场景表结构主键为idchunk key 列设为非主键列pid快照分片Split 01 pid 3Split 13 pid 5并发操作两个不同的 subtask 并发读取 Split 0 与 Split 1此时发生一次更新操作将id0的pid从2改为4且该更新恰好落在两个分片的水位区间内结果Split 0 中含记录[id0, pid2]Split 1 中含记录[id0, pid4]。由于无法保证这两条记录的处理顺序id0最终落库的pid既可能是2也可能是4从而产生数据不一致。因此对于有主键的表务必使用主键列作为 chunk key。DataStream API 使用方式连接器提供两代 DataStream 用法基于SourceFunction的单并行读取旧接口以及基于Db2SourceBuilder的并行增量读取推荐自 3.1.0 起。方式一传统 Db2SourceSourceFunctionimport org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.source.SourceFunction; import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema; public class Db2SourceExample { public static void main(String[] args) throws Exception { SourceFunctionString db2Source Db2Source.Stringbuilder() .hostname(yourHostname) .port(50000) .database(yourDatabaseName) // 设置被捕获的数据库 .tableList(yourSchemaName.yourTableName) // 设置被捕获的表 .username(yourUsername) .password(yourPassword) .deserializer( new JsonDebeziumDeserializationSchema()) // 将 SourceRecord 转换为 JSON 字符串 .build(); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 开启 checkpoint env.enableCheckpointing(3000); env.addSource(db2Source) .print() .setParallelism(1); // sink 使用并行度 1 以保持消息顺序 env.execute(Print Db2 Snapshot Change Stream); } }该入口定义于 Db2Source.java其 Builder 默认port50000、startupOptionsStartupOptions.initial()内部同样以 DebeziumDb2Connector为引擎。方式二并行增量 SourceDb2IncrementalSource推荐自 3.1.0 起可基于Db2SourceBuilder构建支持并行快照与增量读取的Db2IncrementalSourceimport org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.cdc.connectors.base.options.StartupOptions; import org.apache.flink.cdc.connectors.db2.source.Db2SourceBuilder; import org.apache.flink.cdc.connectors.db2.source.Db2SourceBuilder.Db2IncrementalSource; import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema; public class Db2ParallelSourceExample { public static void main(String[] args) throws Exception { Db2IncrementalSourceString db2Source new Db2SourceBuilder() .hostname(localhost) .port(50000) .databaseList(TESTDB) .tableList(DB2INST1.CUSTOMERS) .username(flink) .password(flinkpw) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build(); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 开启 checkpoint env.enableCheckpointing(3000); // 设置 source 并行度为 2 env.fromSource(db2Source, WatermarkStrategy.noWatermarks(), Db2IncrementalSource) .setParallelism(2) .print() .setParallelism(1); env.execute(Print DB2 Snapshot Change Stream); } }Db2IncrementalSource继承自 JDBC 增量读取基类JdbcIncrementalSource其 Builder 提供splitSize、splitMetaGroupSize、distributionFactorUpper/Lower、fetchSize、connectTimeout、connectMaxRetries、connectionPoolSize、chunkKeyColumn、includeSchemaChanges、closeIdleReaders、skipSnapshotBackfill、assignUnboundedChunkFirst等丰富的构建方法见 Db2SourceBuilder.java与 SQL 侧参数一一对应可满足生产级调优需求。可用的监控指标指标系统可帮助了解分片分发与快照读取的进展支持以下 Flink metrics均为 Gauge 类型GroupNameType说明namespace.schema.tableisSnapshottingGauge表是否处于快照读取阶段namespace.schema.tableisStreamReadingGauge表是否处于增量读取阶段namespace.schema.tablenumTablesSnapshottedGauge已完成快照读取的表数量namespace.schema.tablenumTablesRemainingGauge尚未进行快照读取的表数量namespace.schema.tablenumSnapshotSplitsProcessedGauge正在处理的分片数量namespace.schema.tablenumSnapshotSplitsRemainingGauge尚未处理的分片数量namespace.schema.tablenumSnapshotSplitsFinishedGauge已处理完成的分片数量namespace.schema.tablesnapshotStartTimeGauge快照读取阶段开始的时间namespace.schema.tablesnapshotEndTimeGauge快照读取阶段结束的时间注意Group 名称为namespace.schema.table其中namespace是实际的数据库名schema是实际的 schema 名table是实际的表名对于 Db2Group 名称形如test_database.test_schema.test_table。通过这些指标可以实时观察大表快照的进度、剩余分片数以及快照起止时间便于在流式任务运行期间评估数据接入是否落后。数据类型映射Db2 类型与 Flink SQL 类型的映射关系如下表对应 Db2TypeUtils.java 中的转换逻辑Db2 类型Flink SQL 类型备注SMALLINTSMALLINTINTEGERINTBIGINTBIGINTREALFLOATDOUBLEDOUBLENUMERIC(p, s)DECIMAL(p, s)DECIMAL(p, s)精度与标度继承自 Db2 列定义DATEDATETIMETIMETIMESTAMP [(p)]TIMESTAMP [(p)]CHARACTER(n)CHAR(n)VARCHAR(n)VARCHAR(n)BINARY(n)BINARY(n)VARBINARY(n)VARBINARY(n)BLOBCLOBDBCLOBBYTESVARGRAPHICXMLSTRING源码层面的映射细节从 Db2TypeUtils.java 可以看出映射以 JDBCTypes常量为依据CHAR/VARCHAR/SQLXML/CLOB → STRINGBLOB/BINARY/VARBINARY → BYTESFLOAT/REAL → FLOATDECIMAL/NUMERIC → DECIMAL(column.length(), column.scale())可空性会保留若源列非空column.isOptional() falseFlink 类型会被标记为notNull()保证下游类型系统的精确性对未支持的 Db2 类型default分支会抛出UnsupportedOperationException提示“Dont support DB2 type yet”。快速排障要点结合上文内容将常见问题与解决思路整理如下便于快速定位现象可能原因排查/解决方法连接报ClassNotFoundException: com.ibm.db2.jcc.DB2Driver未手动添加 IPLA 协议的 Db2 JDBC 驱动按「依赖引入」章节补充com.ibm.db2.jcc:db2jcc:db2jcc4建表时报 Invalid value for option scan.startup.mode参数取值非法仅支持initial与latest-offset源码见 Db2TableSourceFactory.javascan.startup.mode不生效与debezium.snapshot.mode同时配置二选一使用含BOOLEAN列的表无增量数据SQL Replication 不支持 BOOLEAN改用其他类型或仅接受快照数据无主键表报错未配置 chunk key配置scan.incremental.snapshot.chunk.key-column指定非空列数据不一致主键表误用非主键列作为 chunk key改回主键列或主键第一列作为 chunk key大表快照触发 TaskManager OOM无界 chunk 被最后分配保持scan.incremental.snapshot.unbounded-chunk-first.enabledtrue默认总结Db2 CDC Connector 为 Flink 生态提供了从 IBM Db2 到实时数仓/数据湖的可靠通道initial模式先并行快照、再无缝衔接 redo log 增量latest-offset模式则适合只关心启动后变更的场景。本文已覆盖依赖配置含 IPLA 驱动的手动引入、Flink SQL 建表、16 项连接器参数、虚拟元数据列、DataStream 两种 API、无主键表支持及其一致性语义边界、监控指标与完整类型映射。需要进一步查阅的仓库入口包括Db2 CDC 官方文档、表源工厂、增量 Source Builder 以及 集成测试用例。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考