Oracle CDC 抽取 BLOB 字段完整踩坑与解决指南
2026/7/20 17:39:29 网站建设 项目流程

前些天发现了一个巨牛的人工智能学习网站,通俗易懂,风趣幽默,忍不住分享一下给大家。点击跳转到网站:https://www.captainai.net/dongkelun

前言

之前在 Flink CDC 2 Kafka 中总结过 Flink CDC 对 Blob 字段的支持,测试使用的是MySQL CDC,配置debezium.binary.handling.mode=base64即可将 Blob 图片转 base64 字符串正常同步,当时结论是"Blob 字段支持没问题"。

但后续在Oracle CDC场景中实际操作发现,Blob 字段的支持情况与数据库类型强相关,不能一概而论。Oracle 的 LogMiner 机制对 LOB 类型的处理与 MySQL 截然不同,需要额外的配置和源码级别的修复才能正常工作。本文记录 Oracle CDC 抽取 Blob 字段过程中遇到的全部问题与解决方案。

版本

  • Flink 1.15.3
  • CDC 2.4.2
  • Debezium 1.9.7.Final
  • Oracle 11G / 12C

测试环境与 SQL

公共 Flink SQL 前缀

以下 SQL 在所有场景中通用,放在 SQL 文件最前面(仅作测试示例,实际使用时请按需调整):

setyarn.application.name=cdc_oracle2mysql;setparallelism.default=1;settaskmanager.memory.process.size=3g;settaskmanager.numberOfTaskSlots=1;setexecution.checkpointing.interval=1000;setstate.checkpoints.dir=hdfs:///flink/checkpoints/cdc_oracle2mysql;setexecution.target=yarn-per-job;setexecution.checkpointing.externalized-checkpoint-retention=RETAIN_ON_CANCELLATION;setpipeline.operator-chaining=false;

Flink SQL(Oracle → MySQL)

-- CDC源表CREATETABLEoracle_cdc_source(IDintPRIMARYKEYNOTENFORCED,NAME string,IMG BYTES)WITH('connector'='oracle-cdc',-- 'url' = 'jdbc:oracle:thin:@192.168.1.100:1521:ORCLCDB', -- SID 格式'url'='jdbc:oracle:thin:@//192.168.1.100:1521/ORCLCDB',-- serviceName 格式'hostname'='192.168.1.100','port'='1521','username'='cdc_user','password'='******','database-name'='ORCLCDB','schema-name'='CDC_SCHEMA','table-name'='CDC_SOURCE','debezium.lob.enabled'='true','debezium.log.mining.strategy'='online_catalog','debezium.log.mining.continuous.mine'='true');-- MySQL Sink表CREATETABLEoracle_cdc_sink(IDINTPRIMARYKEYNOTENFORCED,NAME string,IMG BYTES)WITH('connector'='jdbc','driver'='com.mysql.cj.jdbc.Driver','url'='jdbc:mysql://192.168.1.100:3306/test_db?useSSL=false&serverTimezone=Asia/Shanghai&characterEncoding=utf-8','username'='test_user','password'='******','table-name'='CDC_SINK');INSERTINTOoracle_cdc_sinkSELECTID,NAME,IMGFROMoracle_cdc_source;

注:测试时可通过 DBeaver 编辑列 → 从文件导入来插入/更新 BLOB 图片数据。

Flink SQL(Oracle → 达梦)

注:Flink JDBC connector 默认不支持达梦数据库,需要修改源码添加达梦方言支持。

-- 达梦 Sink表CREATETABLEoracle_cdc_sink(IDINTPRIMARYKEYNOTENFORCED,NAME string,IMG BYTES)WITH('connector'='jdbc','driver'='dm.jdbc.driver.DmDriver','url'='jdbc:dm://192.168.1.200:5236','username'='TEST_DBA','password'='******','table-name'='CDC_SINK');INSERTINTOoracle_cdc_sinkSELECTID,NAME,IMGFROMoracle_cdc_source;

第一阶段:不设置 LOB 参数(默认行为)

现象

最开始基于 CDC 2.3.0 进行测试(与 Flink CDC 2 Kafka 测试版本一致),只使用默认参数:

字段结果
NAME(普通字段)正常同步
IMG(BLOB 字段)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.Final,CDC 2.4.2 依赖的是 1.9.7.Final,两个版本对 LOB 支持可能存在差异,因此尝试升级。

Flink 版本为 1.15,CDC 兼容情况如下:

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说明内部传递的startScn0

调试日志分析

为定位根因,在LogMinerStreamingChangeEventSource.java中添加了调试日志(调试源码已提交到 https://gitee.com/dongkelun/flink-cdc.git,分支debug-log-oracle-lob-scn),日志清晰地展示了问题行为:

DEBUG 1 - from offsetContext: startScn=13920606, snapshotScn=null DEBUG 2 - before snapshotScn check: startScn=13920606, snapshotScn=null DEBUG 3 - after compute: startScn=13920606 DEBUG 4 - before startMiningSession: startScn=13920606 DEBUG 5 - after calculateEndScn: startScn=13920606, endScn=13920610 DEBUG 6 - before startMiningSession call: startScn=13920606, endScn=13920610 DEBUG 7 - process returned: startScn=0 ↑ process() 返回了 0 DEBUG 5 - after calculateEndScn: startScn=0, endScn=13920615 DEBUG 6 - before startMiningSession call: startScn=0, endScn=13920615 (循环重复 5 次,startScn 一直为 0) DEBUG startMiningSession exception: startScn=0, endScn=13920650, errorCode=1291

分析

  1. 第 1 轮循环:startScn=13920606正常,processor.process()返回0startScn被改为 0
  2. 第 2~6 轮循环:startScn=0startMiningSession(0+1=1)持续报 ORA-01291。注意endScn在每次循环中持续增长(13920615→…→13920650),说明数据库在正常运行,但startScn一直为 0 导致问题无法自行恢复
  3. 最终在第 6 轮重试次数耗尽,异常抛出

结论:根本问题是processor.process()返回了 0,需要进一步分析processor内部calculateNewStartScn()的代码路径。


第三阶段:根因定位与 flink-cdc 源码修改尝试

尝试过的修复方案

processor.process()返回后尝试修正 SCN:

方案 1:在 process() 返回 0 时保留原 startScn
→ 不报错,但取不到任何数据(BLOB 和其他字段都取不到)

方案 2:在 process() 返回 0 时使用 firstScn
→ 不报错,但同样取不到任何数据

结论

根因在 Debezium 源码中,不在 flink-cdc:

  • AbstractLogMinerEventProcessor.calculateNewStartScn()在 LOB 模式下返回 0
  • 返回 0 后 LOB 数据可被正常抽取(观察到的现象,是否为必要条件不确定)
  • 但 SCN=0 导致 LogMiner 启动 ORA-01291
  • flink-cdc 层面无法绕过——强行修正 SCN 后 LOB 数据抽取逻辑失效

三种场景行为对比

场景报错BLOB 数据其他字段可用性
不设debezium.lob.enablednull正常可用,但 BLOB 为空
lob.enabled=trueORA-01291正常正常不可用(频繁重启)
lob.enabled=true+ flink-cdc 层强制修正 SCNnullnull可用但全字段无数据

第四阶段:修改 Debezium 源码

定位的具体文件与方法

文件方法说明
MemoryLogMinerEventProcessor.javacalculateNewStartScn()AbstractLogMinerEventProcessor的子类,默认实现
AbstractInfinispanLogMinerEventProcessor.javacalculateNewStartScn()AbstractLogMinerEventProcessor的子类,Infinispan 模式

AbstractLogMinerEventProcessor.calculateNewStartScn()仅有以上两个实现类。

旧代码分析

MemoryLogMinerEventProcessor.javacalculateNewStartScn()方法的 LOB 分支旧代码:

if(getConfig().isLobEnabled()){// 情况 A:事务缓存为空 且 maxCommittedScn 有值if(transactionCache.isEmpty()&&!maxCommittedScn.isNull()){// 直接将 offset 设为 maxCommittedScn// 问题:maxCommittedScn 是所有历史会话的最大提交 SCN,首次启动时为 0// 调试日志确认进入此分支(DEBUG LOB PATH A: cache empty, maxCommittedScn=0)// 下一轮从 SCN=0 开始挖掘,LogMiner 启动时加 1 变成 SCN=1 → ORA-01291// 即使 maxCommittedScn 不为 0,也可能是一个陈旧值,// 导致反复挖掘 [maxCommittedScn+1, endScn] 区间 → 死循环offsetContext.setScn(maxCommittedScn);dispatcher.dispatchHeartbeatEvent(partition,offsetContext);}else{// 情况 B:事务缓存非空,或 maxCommittedScn 为 nullabandonTransactions(getConfig().getLogMiningTransactionRetention());finalScnminStartScn=getTransactionCacheMinimumScn();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-DskipTests=true-Dspotless.check.skip=true-Dcheckstyle.skip=true
DEBUG LOB PATH A: cache empty, maxCommittedScn=0 DEBUG LOB RETURN: 0

日志确认进入了情况 A(事务缓存为空且maxCommittedScn=0),直接将 offset SCN 设为 0 并返回。

实际上通过源码分析也可以得出结论:情况 B1 将 SCN 设为minStartScn - 1(正常 SCN),情况 B 中minStartScn为 null 时什么都不做(直接到共同的return),均不会产生 0。返回 0 只有情况 A 中maxCommittedScn=0这一种可能。调试日志的作用是验证了实际运行时确实进入了情况 A,排除了代码审查中的不确定性。

新代码

if(getConfig().isLobEnabled()){// 清理过期事务和缓存(移到条件分支外,保证始终执行)abandonTransactions(getConfig().getLogMiningTransactionRetention());finalScnminStartScn=getTransactionCacheMinimumScn();if(!minStartScn.isNull()){recentlyProcessedTransactionsCache.entrySet().removeIf(...);schemaChangesCache.removeIf(...);}// LGWR buffer 未完全落盘时,用 lastProcessedScn 修正 endScn// 使 LOB 分支与非 LOB 分支的 LGWR 处理保持一致if(!getLastProcessedScn().isNull()&&getLastProcessedScn().compareTo(endScn)<0){endScn=getLastProcessedScn();}// 统一使用 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
可能返回 0maxCommittedScn=0且 cache 空时不会无此问题

解决的问题

  1. 修复挖掘死循环:旧代码情况 A(cache 空且maxCommittedScn非 null)将 offset 设为maxCommittedScn,该值来自所有历史会话的 CommitScn,并非当前 batch 的实际进度。当maxCommittedScn陈旧时(如首次启动为 0),每次 process 返回同一个陈旧值,下一轮又从该值开始挖掘,反复 mining 同一区间。
  2. 修复 ORA-01291:旧代码情况 B 回退到minStartScn - 1,该 SCN 可能过于陈旧,对应的 redo/archive log 已被清理,LogMiner 启动时因找不到日志文件而报错。新代码统一用endScn推进,不回退。
  3. 避免重复挖掘 LOB 事件:LOB 操作(SELECT_LOB_LOCATORLOB_WRITELOB_ERASE)已在当前 SCN 范围被捕获并存于事务缓存,回退重放是多余的。
  4. 兜底 endScn 超前:当lastProcessedScn < endScn时,用lastProcessedScn修正endScn,避免挖到尚未完全落盘的 redo 数据。

验证

  • MemoryLogMinerEventProcessor.java已验证通过,ORA-01291 不再出现
  • AbstractInfinispanLogMinerEventProcessor.java⚠️未验证(Flink Oracle CDC 默认log.mining.buffer.type=memory,不走 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.IMG=OLD.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(),统一使用endScn
2. 目标端创建触发器__debezium_unavailable_value覆盖真实 BLOB 数据BEFORE UPDATE 触发器拦截占位符

最佳实践参数配置

'debezium.lob.enabled'='true','debezium.log.mining.strategy'='online_catalog','debezium.log.mining.continuous.mine'='true'

注意事项

  • 不修改 BLOB 的 UPDATE:即使修复了 ORA-01291,__debezium_unavailable_value问题依然存在,必须配合触发器使用
  • Infinispan 模式:本次修改也覆盖了 Infinispan 实现,但未实测验证

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询