Oracle CDC 抽取 BLOB 字段完整踩坑与解决指南
前些天发现了一个巨牛的人工智能学习网站通俗易懂风趣幽默忍不住分享一下给大家。点击跳转到网站https://www.captainai.net/dongkelun前言之前在 Flink CDC 2 Kafka 中总结过 Flink CDC 对 Blob 字段的支持测试使用的是MySQL CDC配置debezium.binary.handling.modebase64即可将 Blob 图片转 base64 字符串正常同步当时结论是Blob 字段支持没问题。但后续在Oracle CDC场景中实际操作发现Blob 字段的支持情况与数据库类型强相关不能一概而论。Oracle 的 LogMiner 机制对 LOB 类型的处理与 MySQL 截然不同需要额外的配置和源码级别的修复才能正常工作。本文记录 Oracle CDC 抽取 Blob 字段过程中遇到的全部问题与解决方案。版本Flink 1.15.3CDC 2.4.2Debezium 1.9.7.FinalOracle 11G / 12C测试环境与 SQL公共 Flink SQL 前缀以下 SQL 在所有场景中通用放在 SQL 文件最前面仅作测试示例实际使用时请按需调整setyarn.application.namecdc_oracle2mysql;setparallelism.default1;settaskmanager.memory.process.size3g;settaskmanager.numberOfTaskSlots1;setexecution.checkpointing.interval1000;setstate.checkpoints.dirhdfs:///flink/checkpoints/cdc_oracle2mysql;setexecution.targetyarn-per-job;setexecution.checkpointing.externalized-checkpoint-retentionRETAIN_ON_CANCELLATION;setpipeline.operator-chainingfalse;Flink SQLOracle → MySQL-- CDC源表CREATETABLEoracle_cdc_source(IDintPRIMARYKEYNOTENFORCED,NAME string,IMG BYTES)WITH(connectororacle-cdc,-- url jdbc:oracle:thin:192.168.1.100:1521:ORCLCDB, -- SID 格式urljdbc:oracle:thin://192.168.1.100:1521/ORCLCDB,-- serviceName 格式hostname192.168.1.100,port1521,usernamecdc_user,password******,database-nameORCLCDB,schema-nameCDC_SCHEMA,table-nameCDC_SOURCE,debezium.lob.enabledtrue,debezium.log.mining.strategyonline_catalog,debezium.log.mining.continuous.minetrue);-- MySQL Sink表CREATETABLEoracle_cdc_sink(IDINTPRIMARYKEYNOTENFORCED,NAME string,IMG BYTES)WITH(connectorjdbc,drivercom.mysql.cj.jdbc.Driver,urljdbc:mysql://192.168.1.100:3306/test_db?useSSLfalseserverTimezoneAsia/ShanghaicharacterEncodingutf-8,usernametest_user,password******,table-nameCDC_SINK);INSERTINTOoracle_cdc_sinkSELECTID,NAME,IMGFROMoracle_cdc_source;注测试时可通过 DBeaver 编辑列 → 从文件导入来插入/更新 BLOB 图片数据。Flink SQLOracle → 达梦注Flink JDBC connector 默认不支持达梦数据库需要修改源码添加达梦方言支持。-- 达梦 Sink表CREATETABLEoracle_cdc_sink(IDINTPRIMARYKEYNOTENFORCED,NAME string,IMG BYTES)WITH(connectorjdbc,driverdm.jdbc.driver.DmDriver,urljdbc:dm://192.168.1.200:5236,usernameTEST_DBA,password******,table-nameCDC_SINK);INSERTINTOoracle_cdc_sinkSELECTID,NAME,IMGFROMoracle_cdc_source;第一阶段不设置 LOB 参数默认行为现象最开始基于 CDC 2.3.0 进行测试与 Flink CDC 2 Kafka 测试版本一致只使用默认参数字段结果NAME普通字段正常同步IMGBLOB 字段null尝试排查在 CDC 2.3.0 上添加debezium.lob.enabled true参数BLOB 字段依然为 null猜测 2.3.0 的 Oracle connector 对 LOB 支持不完整尝试升级版本。版本升级CDC 2.3.0 依赖的 Debezium 版本是 1.6.4.FinalCDC 2.4.2 依赖的是 1.9.7.Final两个版本对 LOB 支持可能存在差异因此尝试升级。Flink 版本为 1.15CDC 兼容情况如下CDC 版本说明2.4.2兼容 Flink 1.15 的最高稳定版本3.0也兼容 1.15但作为 3.x 的首个版本可能不够稳定未经过充分验证经实测CDC 3.0 添加debezium.lob.enabled true后的表现与 2.4.2 完全一致都会报 ORA-01291也都能取到数据。考虑到稳定性最终采用CDC 2.4.2。升级到 2.4.2 后添加debezium.lob.enabled true参数BLOB 字段终于不再为 null——但紧接着遇到了新问题。第二阶段启用 LOB 参数ORA-01291 问题现象添加debezium.lob.enabled true后数据包括 BLOB 字段能正常抽取但持续报ORA-01291: missing logfile。由于持续报错连接器频繁触发重启导致以下影响日志量激增LogMiner 会话不断创建、泄漏连接器可能陷入无限重启循环不可用错误日志Caused by: java.sql.SQLException: ORA-01291: missing logfile ORA-06512: at SYS.DBMS_LOGMNR, line 58 ORA-06512: at line 1 at oracle.jdbc.driver.T4CTTIoer11.processError(T4CTTIoer11.java:509) at oracle.jdbc.driver.T4CTTIoer11.processError(T4CTTIoer11.java:461) at oracle.jdbc.driver.T4C8Oall.processError(T4C8Oall.java:1104) at oracle.jdbc.driver.T4CTTIfun.receive(T4CTTIfun.java:550) at oracle.jdbc.driver.T4CTTIfun.doRPC(T4CTTIfun.java:268) at oracle.jdbc.driver.T4C8Oall.doOALL(T4C8Oall.java:655) at oracle.jdbc.driver.T4CStatement.doOall8(T4CStatement.java:229) at oracle.jdbc.driver.T4CStatement.doOall8(T4CStatement.java:41) at oracle.jdbc.driver.T4CStatement.executeForRows(T4CStatement.java:928) at oracle.jdbc.driver.OracleStatement.doExecuteWithTimeout(OracleStatement.java:1205) at oracle.jdbc.driver.OracleStatement.executeInternal(OracleStatement.java:1823) at oracle.jdbc.driver.OracleStatement.execute(OracleStatement.java:1778) at oracle.jdbc.driver.OracleStatementWrapper.execute(OracleStatementWrapper.java:303) at io.debezium.jdbc.JdbcConnection.executeWithoutCommitting(JdbcConnection.java:1446) at io.debezium.connector.oracle.logminer.LogMinerStreamingChangeEventSource.startMiningSession(LogMinerStreamingChangeEventSource.java:679) at io.debezium.connector.oracle.logminer.LogMinerStreamingChangeEventSource.execute(LogMinerStreamingChangeEventSource.java:242) ... 8 more Caused by: Error : 1291, Position : 0, Sql BEGIN sys.dbms_logmnr.start_logmnr(startScn 1, endScn 13223261, OPTIONS DBMS_LOGMNR.DICT_FROM_ONLINE_CATALOG DBMS_LOGMNR.CONTINUOUS_MINE DBMS_LOGMNR.NO_ROWID_IN_STMT);END;, OriginalSql BEGIN sys.dbms_logmnr.start_logmnr(startScn 1, endScn 13223261, OPTIONS DBMS_LOGMNR.DICT_FROM_ONLINE_CATALOG DBMS_LOGMNR.CONTINUOUS_MINE DBMS_LOGMNR.NO_ROWID_IN_STMT);END;, Error Msg ORA-01291: missing logfile ORA-06512: at SYS.DBMS_LOGMNR, line 58 ORA-06512: at line 1 at oracle.jdbc.driver.T4CTTIoer11.processError(T4CTTIoer11.java:513) ... 23 more通过完整的堆栈可以看到异常由LogMinerStreamingChangeEventSource.startMiningSession()方法抛出该方法由execute()方法调用。startMiningSession的参数startScn来自调用方需要结合调试日志进一步定位startScn的来源。关键信息startScn 1。查看源码startMiningSession()发现调用startLogMinerStatement()时会在传入的startScn上加 1因此日志中的1说明内部传递的startScn是0。调试日志分析为定位根因在LogMinerStreamingChangeEventSource.java中添加了调试日志调试源码已提交到 https://gitee.com/dongkelun/flink-cdc.git分支debug-log-oracle-lob-scn日志清晰地展示了问题行为DEBUG 1 - from offsetContext: startScn13920606, snapshotScnnull DEBUG 2 - before snapshotScn check: startScn13920606, snapshotScnnull DEBUG 3 - after compute: startScn13920606 DEBUG 4 - before startMiningSession: startScn13920606 DEBUG 5 - after calculateEndScn: startScn13920606, endScn13920610 DEBUG 6 - before startMiningSession call: startScn13920606, endScn13920610 DEBUG 7 - process returned: startScn0 ↑ process() 返回了 0 DEBUG 5 - after calculateEndScn: startScn0, endScn13920615 DEBUG 6 - before startMiningSession call: startScn0, endScn13920615 循环重复 5 次startScn 一直为 0 DEBUG startMiningSession exception: startScn0, endScn13920650, errorCode1291分析第 1 轮循环startScn13920606正常processor.process()返回0→startScn被改为 0第 26 轮循环startScn0startMiningSession(011)持续报 ORA-01291。注意endScn在每次循环中持续增长13920615→…→13920650说明数据库在正常运行但startScn一直为 0 导致问题无法自行恢复最终在第 6 轮重试次数耗尽异常抛出结论根本问题是processor.process()返回了 0需要进一步分析processor内部calculateNewStartScn()的代码路径。第三阶段根因定位与 flink-cdc 源码修改尝试尝试过的修复方案在processor.process()返回后尝试修正 SCN方案 1在 process() 返回 0 时保留原 startScn→ 不报错但取不到任何数据BLOB 和其他字段都取不到方案 2在 process() 返回 0 时使用 firstScn→ 不报错但同样取不到任何数据结论根因在 Debezium 源码中不在 flink-cdcAbstractLogMinerEventProcessor.calculateNewStartScn()在 LOB 模式下返回 0返回 0 后 LOB 数据可被正常抽取观察到的现象是否为必要条件不确定但 SCN0 导致 LogMiner 启动 ORA-01291flink-cdc 层面无法绕过——强行修正 SCN 后 LOB 数据抽取逻辑失效三种场景行为对比场景报错BLOB 数据其他字段可用性不设debezium.lob.enabled否null正常可用但 BLOB 为空设lob.enabledtrueORA-01291正常正常不可用频繁重启设lob.enabledtrue flink-cdc 层强制修正 SCN否nullnull可用但全字段无数据第四阶段修改 Debezium 源码定位的具体文件与方法文件方法说明MemoryLogMinerEventProcessor.javacalculateNewStartScn()AbstractLogMinerEventProcessor的子类默认实现AbstractInfinispanLogMinerEventProcessor.javacalculateNewStartScn()AbstractLogMinerEventProcessor的子类Infinispan 模式AbstractLogMinerEventProcessor.calculateNewStartScn()仅有以上两个实现类。旧代码分析MemoryLogMinerEventProcessor.java中calculateNewStartScn()方法的 LOB 分支旧代码if(getConfig().isLobEnabled()){// 情况 A事务缓存为空 且 maxCommittedScn 有值if(transactionCache.isEmpty()!maxCommittedScn.isNull()){// 直接将 offset 设为 maxCommittedScn// 问题maxCommittedScn 是所有历史会话的最大提交 SCN首次启动时为 0// 调试日志确认进入此分支DEBUG LOB PATH A: cache empty, maxCommittedScn0// 下一轮从 SCN0 开始挖掘LogMiner 启动时加 1 变成 SCN1 → ORA-01291// 即使 maxCommittedScn 不为 0也可能是一个陈旧值// 导致反复挖掘 [maxCommittedScn1, endScn] 区间 → 死循环offsetContext.setScn(maxCommittedScn);dispatcher.dispatchHeartbeatEvent(partition,offsetContext);}else{// 情况 B事务缓存非空或 maxCommittedScn 为 nullabandonTransactions(getConfig().getLogMiningTransactionRetention());finalScnminStartScngetTransactionCacheMinimumScn();if(!minStartScn.isNull()){// 回退到 minStartScn - 1// 问题回退到旧的 SCN 时对应的 redo/archive log 可能已被清理// → ORA-01291// 另外LOB 事件SELECT_LOB_LOCATOR, LOB_WRITE, LOB_ERASE// 已在当前 SCN 范围被 LogMiner 捕获并存于事务缓存// 回退重放是多余的recentlyProcessedTransactionsCache.entrySet().removeIf(...);schemaChangesCache.removeIf(...);offsetContext.setScn(minStartScn.subtract(Scn.valueOf(1)));dispatcher.dispatchHeartbeatEvent(partition,offsetContext);}}returnoffsetContext.getScn();}旧代码调试分析为了确认calculateNewStartScn()中具体哪个代码路径返回了 0在 Debezium 旧代码的 LOB 分支中增加了路径日志源码已提交到 https://gitee.com/dongkelun/debezium.git分支debug-old-lob-path运行后输出如下调试代码未通过代码格式检查编译时可跳过mvn cleaninstall-DskipTeststrue-Dspotless.check.skiptrue-Dcheckstyle.skiptrueDEBUG LOB PATH A: cache empty, maxCommittedScn0 DEBUG LOB RETURN: 0日志确认进入了情况 A事务缓存为空且maxCommittedScn0直接将 offset SCN 设为 0 并返回。实际上通过源码分析也可以得出结论情况 B1 将 SCN 设为minStartScn - 1正常 SCN情况 B 中minStartScn为 null 时什么都不做直接到共同的return均不会产生 0。返回 0 只有情况 A 中maxCommittedScn0这一种可能。调试日志的作用是验证了实际运行时确实进入了情况 A排除了代码审查中的不确定性。新代码if(getConfig().isLobEnabled()){// 清理过期事务和缓存移到条件分支外保证始终执行abandonTransactions(getConfig().getLogMiningTransactionRetention());finalScnminStartScngetTransactionCacheMinimumScn();if(!minStartScn.isNull()){recentlyProcessedTransactionsCache.entrySet().removeIf(...);schemaChangesCache.removeIf(...);}// LGWR buffer 未完全落盘时用 lastProcessedScn 修正 endScn// 使 LOB 分支与非 LOB 分支的 LGWR 处理保持一致if(!getLastProcessedScn().isNull()getLastProcessedScn().compareTo(endScn)0){endScngetLastProcessedScn();}// 统一使用 endScn不再回退到 minStartScn - 1 或使用 maxCommittedScn// 旧逻辑中 maxCommittedScn 可能陈旧导致死循环回退旧 SCN 会导致 ORA-01291// LOB 事件已在当前 SCN 范围被捕获无须回退重放offsetContext.setScn(endScn);dispatcher.dispatchHeartbeatEvent(partition,offsetContext);returnoffsetContext.getScn();}完整源码已提交到 https://gitee.com/dongkelun/debezium.git分支v1.9.7.Final-fix-lob-scn-calc。修改前后对比对比维度LOB 分支修改前LOB 分支修改后非 LOB 分支未改动LGWR 处理无lastProcessedScn修正新增与非 LOB 一致有lastProcessedScn修正缓存清理else分支中执行移到外层无条件执行else分支中执行SCN 推进策略cache 空用maxCommittedScn否则minStartScn - 1统一用endScncache 空用endScn否则minStartScn - 1可能返回 0当maxCommittedScn0且 cache 空时不会无此问题解决的问题修复挖掘死循环旧代码情况 Acache 空且maxCommittedScn非 null将 offset 设为maxCommittedScn该值来自所有历史会话的 CommitScn并非当前 batch 的实际进度。当maxCommittedScn陈旧时如首次启动为 0每次 process 返回同一个陈旧值下一轮又从该值开始挖掘反复 mining 同一区间。修复 ORA-01291旧代码情况 B 回退到minStartScn - 1该 SCN 可能过于陈旧对应的 redo/archive log 已被清理LogMiner 启动时因找不到日志文件而报错。新代码统一用endScn推进不回退。避免重复挖掘 LOB 事件LOB 操作SELECT_LOB_LOCATOR、LOB_WRITE、LOB_ERASE已在当前 SCN 范围被捕获并存于事务缓存回退重放是多余的。兜底 endScn 超前当lastProcessedScn endScn时用lastProcessedScn修正endScn避免挖到尚未完全落盘的 redo 数据。验证MemoryLogMinerEventProcessor.java✅已验证通过ORA-01291 不再出现AbstractInfinispanLogMinerEventProcessor.java⚠️未验证Flink Oracle CDC 默认log.mining.buffer.typememory不走 Infinispan 逻辑为保持一致性一并修改__debezium_unavailable_value问题问题描述无论是否修改 Debezium 源码当 BLOB 字段存储的是图片等二进制数据并且执行 UPDATE 操作时仅在不修改 BLOB 字段本身的情况下目标表中的 BLOB 字段会被更新为__debezium_unavailable_value占位符。如果 UPDATE 同时修改了 BLOB 字段本身则数据正常。表现执行 UPDATE 时目标表中的 BLOB 字段被更新为__debezium_unavailable_value覆盖原有的真实数据。具体原因未深入分析可能是 LogMiner 在捕获 LOB 事件时未携带实际内容Debezium 用占位符填充。注意该问题仅在UPDATE时出现INSERT 和 DELETE 均正常。另外如果 BLOB 字段存储的是普通字符串而非图片等二进制数据由于字符串在 CDC 传输过程中会被完整捕获UPDATE 时也不会出现此问题。解决方案触发器拦截MySQL 目标端触发器CREATETRIGGERtrg_block_debezium_placeholder_blob BEFOREUPDATEONtest_db.CDC_SINKFOR EACH ROWBEGINIFCONVERT(NEW.IMGUSINGutf8mb4)LIKE%__debezium_unavailable_value%THENSETNEW.IMGOLD.IMG;ENDIF;END达梦目标端触发器CREATEORREPLACETRIGGERTRG_BLOCK_DEBEZIUM_PLACEHOLDER_BLOB BEFOREUPDATEONCDC_SINKFOR EACH ROWDECLAREV_PLACEHOLDER RAW(100);BEGINV_PLACEHOLDER :UTL_RAW.CAST_TO_RAW(__debezium_unavailable_value);IFDBMS_LOB.INSTR(:NEW.IMG,V_PLACEHOLDER,1,1)0THEN:NEW.IMG ::OLD.IMG;ENDIF;END;触发器原理在 UPDATE 触发前检查新传入的 BLOB 值是否包含__debezium_unavailable_value占位符。如果包含说明上游发送的是占位符而非真实数据将新值还原为数据库中的老值拒绝被覆盖。触发器方案只是临时绕过不是根本解决。根本解决可能需要修改 Debezium 源码让 LogMiner 在捕获 LOB 事件时携带实际数据而非占位符但本文暂不涉及。触发器方案只适用于直接写入数据库的场景。如果通过 Kafka 消费 CDC 数据需要在消费端判断 BLOB 字段值是否等于__debezium_unavailable_value如果是则忽略该字段保留原值。本次需求场景正好是Oracle → 达梦直接写入因此触发器方案可以落地。最终方案总结两步解决步骤解决什么问题方案1. 修改 Debezium 源码ORA-01291 LogMiner 死循环修改calculateNewStartScn()统一使用endScn2. 目标端创建触发器__debezium_unavailable_value覆盖真实 BLOB 数据BEFORE UPDATE 触发器拦截占位符最佳实践参数配置debezium.lob.enabledtrue,debezium.log.mining.strategyonline_catalog,debezium.log.mining.continuous.minetrue注意事项不修改 BLOB 的 UPDATE即使修复了 ORA-01291__debezium_unavailable_value问题依然存在必须配合触发器使用Infinispan 模式本次修改也覆盖了 Infinispan 实现但未实测验证