Skip to content

[oracle] Support row_id metadata column for Oracle CDC (with headers-loss bug fix) - #4486

Open
fightBoxing wants to merge 1 commit into
apache:release-2.4from
fightBoxing:test/oracle-release-2.4
Open

[oracle] Support row_id metadata column for Oracle CDC (with headers-loss bug fix)#4486
fightBoxing wants to merge 1 commit into
apache:release-2.4from
fightBoxing:test/oracle-release-2.4

Conversation

@fightBoxing

Copy link
Copy Markdown

概述

为 Oracle CDC 连接器新增 row_id 元数据列支持,用户可通过 METADATA FROM 'row_id' VIRTUAL 在 Flink SQL 中获取 Oracle 表的 ROWID 伪列。

动机

业务场景需要基于 Oracle ROWID 进行:

  • 数据溯源与精确定位物理行
  • 幂等写入与去重
  • 兼容基于 ROWID 的下游处理

使用方式

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 tabSELECT 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 参构造函数保留 timestampheaders

此修复影响范围:不仅解决 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 值

Add support for extracting Oracle ROWID pseudo-column as a metadata
column via 'METADATA FROM row_id VIRTUAL' syntax in Flink SQL.

Changes:
1. OracleReadableMetaData: Add ROW_ID metadata enum that reads ROWID
   from SourceRecord headers.

2. OracleScanFetchTask (snapshot phase):
   - Rewrite SELECT SQL to append ROWID pseudo-column at the tail:
     'SELECT * FROM tab' -> 'SELECT T0.*, ROWID FROM tab T0'.
   - Manually construct ColumnArray excluding ROWID to bypass Debezium
     ColumnUtils strict validation.
   - Extract ROWID via OracleResultSet.getROWID(N+1) and inject into
     SourceRecord headers by overriding
     SnapshotChangeRecordEmitter#getEmitConnectHeaders.

3. JdbcSourceFetchTaskContext: Fix headers loss bug in
   formatMessageTimestamp(). The previous 8-arg SourceRecord
   constructor dropped headers when rewriting snapshot records; switch
   to the 10-arg constructor to preserve timestamp and headers. This
   fix benefits all JDBC-based Flink CDC connectors when they rely on
   headers for downstream metadata.

Requirements:
- scan.incremental.snapshot.enabled=true (Debezium native snapshot
  path is not covered by this change).
- Table supplemental log enabled for streaming ROWID capture.

Verified with Oracle 12c Standard Edition (12.1.0.2.0). Sample output:
  +I[2, bbb, AABDuHAAGAAAAF3AAA]
  +I[3, ccc, AABDuHAAGAAAAF0AAA]
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant