fightBoxing opened a new pull request, #4486:
URL: https://github.com/apache/flink-cdc/pull/4486

   ## 概述
   
   为 Oracle CDC 连接器新增 `row_id` 元数据列支持,用户可通过 `METADATA FROM 'row_id' VIRTUAL` 在 
Flink SQL 中获取 Oracle 表的 ROWID 伪列。
   
   ## 动机
   
   业务场景需要基于 Oracle ROWID 进行:
   - 数据溯源与精确定位物理行
   - 幂等写入与去重
   - 兼容基于 ROWID 的下游处理
   
   ## 使用方式
   
   ```sql
   CREATE TABLE oracle_source (
       ID INT,
       NAME STRING,
       row_id STRING METADATA FROM 'row_id' VIRTUAL   -- 新增
   ) WITH (
       'connector' = 'oracle-cdc',
       'scan.incremental.snapshot.enabled' = 'true',
       -- ... 其他连接配置
   );
   ```
   
   输出示例:
   ```
   +I[2, bbb, AABDuHAAGAAAAF3AAA]
   +I[3, ccc, AABDuHAAGAAAAF0AAA]
   -D[2, bbb, AABDuHAAGAAAAF3AAA]
   ```
   
   ## 改动清单
   
   ### 1. `OracleReadableMetaData.java` — 新增 ROW_ID 元数据枚举
   从 `SourceRecord.headers()` 中读取 key = `ROWID` 的 header 值。
   
   ### 2. `OracleScanFetchTask.java` — 快照阶段实现
   - **SQL 重写**:`SELECT * FROM tab` → `SELECT T0.*, ROWID FROM tab T0`(ROWID 
放末尾,保持物理列位置 1..N 不变)
   - **绕过 Debezium 严格校验**:手动构造 `ColumnArray`,只包含表原有列,避免 `ColumnUtils.toArray` 因 
ROWID 列不在 schema 中而报错
   - **独立提取 ROWID**:`((OracleResultSet) rs).getROWID(N+1)`
   - **注入 headers**:覆写匿名 `SnapshotChangeRecordEmitter#getEmitConnectHeaders()`
   
   ### 3. `JdbcSourceFetchTaskContext.java` — 关键 bug 修复 ⭐
   `formatMessageTimestamp()` 在快照阶段重构 SourceRecord 时使用了 8 参构造函数,**丢弃了原始 record 
的 headers**。改为 10 参构造函数保留 `timestamp` 和 `headers`。
   
   **此修复影响范围**:不仅解决 Oracle row_id,也惠及所有基于 JDBC 的 Flink CDC 
连接器(MySQL、PostgreSQL、SQL Server、Db2),使它们的增量快照模式下 SourceRecord headers 
得以正确保留,为后续实现其他 headers 相关的元数据列(SCN、Position 等)打好基础。
   
   ### 4. `docs/ROWID_METADATA_IMPLEMENTATION.md` — 方案说明文档
   
   ## 前置要求
   
   - `scan.incremental.snapshot.enabled = true`(本 PR 未修改 Debezium 原生快照路径)
   - 表已开启补充日志:`ALTER TABLE <schema>.<table> ADD SUPPLEMENTAL LOG DATA (ALL) 
COLUMNS`
   - 账户具备 Oracle LogMiner 相关权限
   
   ## 验证
   
   - 编译:`mvn clean package -Dflink.version=1.16.3` ✅
   - 环境:Oracle Database 12c Standard Edition Release 12.1.0.2.0
   - 验证覆盖:快照阶段(INSERT)、流式阶段(INSERT/UPDATE/DELETE)ROWID 均正确输出
   
   ## 数据流
   
   ```
   OracleScanFetchTask (rs.getROWID)
       → getEmitConnectHeaders() 注入 ROWID
           → BufferingSnapshotChangeRecordReceiver 创建 SourceRecord(headers)
               → ChangeEventQueue
                   → IncrementalSourceScanFetcher.pollSplitRecords
                       → JdbcSourceFetchTaskContext.formatMessageTimestamp ⭐ 
(修复保留 headers)
                           → 
OracleReadableMetaData.ROW_ID.read(record.headers())
                               → 输出 row_id 值
   ```
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to