[
https://issues.apache.org/jira/browse/FLINK-40201?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Mohammad Suhel Ksolves Cassandra Experties updated FLINK-40201:
---------------------------------------------------------------
Description:
Here it is — copy and paste directly:
----
h2. Background
Flink CDC 3.6's Oracle connector vendors several Debezium internal classes that
were written against an older Debezium API. Oracle 26ai (23.26.x) requires
features only available in Debezium 3.4.2.Final (specifically
UnbufferedLogMinerStreamingChangeEventSource introduced in Debezium 2.x). When
Debezium 3.4.2.Final is used, the vendored classes produce runtime errors
because they still carry old API signatures. Simply swapping the Debezium JAR
is not sufficient — several source-level fixes are required to align the
vendored classes with the Debezium 3.4.2 API.
h2. Problem
Oracle CDC connector in Flink CDC 3.6 fails with Oracle 26ai (CDB/PDB mode)
when using Debezium 3.4.2.Final. Multiple runtime errors prevent snapshot and
streaming CDC from working end-to-end:
# {*}NoSuchMethodError on LogMinerChangeRecordEmitter{*}: The vendored
LogMinerChangeRecordEmitter carries a 10-arg constructor (from an older
Debezium API) but Debezium 3.4.2's UnbufferedLogMinerStreamingChangeEventSource
calls a 9-arg constructor (rowId parameter was removed in Debezium 2.x).
# {*}NullPointerException in ResumePositionProvider.close(){*}: The close()
method calls sessionContext.close() without null-guarding sessionContext. If
lazy initialisation fails after connection is assigned but before
sessionContext is assigned, close() produces a secondary NPE that hides the
root cause exception.
# {*}Invalid JDBC connection string in ResumePositionProvider{*}: When the
connector config carries "database.*" prefixed keys, OracleConnection receives
keys like "database.hostname" instead of "hostname", producing a malformed JDBC
URL and connection failure.
# {*}ORA-01031 on row count query{*}: OracleUtils.queryApproximateRowCnt()
used ANALYZE TABLE which requires the ANALYZE ANY system privilege. Standard
CDC users do not have this privilege. Fixed by querying ALL_TABLES.NUM_ROWS
instead.
# {*}ORA-00942 / schema recovery failure{*}: readSchemaForCapturedTables()
fails in CDB/PDB mode because Oracle JDBC returns the PDB name as TABLE_CAT in
JDBC metadata results, but the TableId has a null catalog. Fixed by using a
catalog-agnostic filter when calling readSchema().
# {*}NullPointerException on absent ROWID Connect header{*}: Debezium 3.4.2
removed the getEmitConnectHeaders() override that injected a ROWID into Kafka
Connect headers. OracleSourceFetchTaskContext.isRecordBetween() then crashes on
tables without a primary key because headers.iterator().next() throws
NoSuchElementException when no header is present.
# {*}ClassNotFoundException at runtime{*}: debezium-common classes (e.g.
RelationalDatabaseSchema) are referenced at runtime but were not included in
the Maven shade assembly, causing ClassNotFoundException when the fat JAR is
used.
h2. Steps to Reproduce
# Set up Oracle 26ai in CDB/PDB mode (e.g. FREE CDB + FREEPDB1)
# Upgrade Debezium version in pom.xml to 3.4.2.Final
# Run Oracle CDC connector with scan.startup.mode=initial
# Observe failures listed above during snapshot and streaming phases
h2. Fix
* Update vendored LogMinerChangeRecordEmitter to 9-arg constructor matching
Debezium 3.4.2
* Add local ResumePositionProvider override: null-guard on sessionContext in
close() and catch block; strip "database." prefix from JdbcConfiguration before
creating OracleConnection
* Fix OracleUtils.queryApproximateRowCnt() to use ALL_TABLES.NUM_ROWS instead
of ANALYZE TABLE
* Fix OracleSchema.getTableSchema() to use a catalog-agnostic TableFilter for
readSchema() so it works correctly in Oracle CDB/PDB mode
* Guard isRecordBetween() against absent ROWID Connect header; fall back to
SCN-range inclusion for tables without a primary key (multi-chunk PK-less
tables should set scan.incremental.snapshot.chunk.key-column explicitly)
* Include debezium-common in Maven shade assembly to prevent
ClassNotFoundException at runtime
* Add EmbeddedFlinkSchemaHistory and RelationalDatabaseSchema adapters for
Debezium 3.x API
h2. Verification
Tested end-to-end pipeline: Oracle 26ai FREE (CDB/FREEPDB1) → Flink CDC
3.6-SNAPSHOT → StarRocks sink. Confirmed 3 snapshot rows and a live CDC UPDATE
propagated correctly using Debezium 3.4.2.Final and scan.startup.mode=initial.
h2. Environment
* Flink CDC: 3.6-SNAPSHOT
* Debezium: 3.4.2.Final
* Oracle: 26ai (FREE edition, CDB + FREEPDB1 PDB)
* Java: 17
was:
*Background*
Flink CDC 3.6's Oracle connector vendors several Debezium internal classes
that were
written against an older Debezium API. Oracle 26ai (23.26.x) requires
features only
available in Debezium 3.4.2.Final (specifically
UnbufferedLogMinerStreamingChangeEventSource
introduced in Debezium 2.x). When Debezium 3.4.2.Final is used, the vendored
classes
produce runtime errors because they still carry old API signatures. Simply
swapping
the Debezium JAR is not sufficient — several source-level fixes are required
to align
the vendored classes with the Debezium 3.4.2 API.
*Problem*
Oracle CDC connector in Flink CDC 3.6 fails with Oracle 26ai (CDB/PDB mode)
when
using Debezium 3.4.2.Final. Multiple runtime errors prevent snapshot and
streaming
CDC from working end-to-end:
1. {*}{{*}}NoSuchMethodError on LogMinerChangeRecordEmitter{{*}}{*}: The
vendored
LogMinerChangeRecordEmitter carries a 10-arg constructor (from an older
Debezium API)
but Debezium 3.4.2's UnbufferedLogMinerStreamingChangeEventSource calls a
9-arg
constructor (rowId parameter was removed in Debezium 2.x).
2. {*}{{*}}NullPointerException in ResumePositionProvider.close(){{*}}{*}:
The close() method
calls sessionContext.close() without null-guarding sessionContext. If lazy
initialisation fails after connection is assigned but before
sessionContext is
assigned, close() produces a secondary NPE that hides the root cause
exception.
3. {*}{{*}}Invalid JDBC connection string in ResumePositionProvider{{*}}{*}:
When the connector
"database.hostname" instead of "hostname", producing a malformed JDBC URL
and
connection failure.
4. {*}{{*}}ORA-01031 on row count query{{*}}{*}:
OracleUtils.queryApproximateRowCnt() used
ANALYZE TABLE which requires the ANALYZE ANY system privilege. Standard CDC
users do not have this privilege. Fixed by querying ALL_TABLES.NUM_ROWS
instead.
5. {*}{{*}}ORA-00942 / schema recovery failure{{*}}{*}:
readSchemaForCapturedTables() fails in
CDB/PDB mode because Oracle JDBC returns the PDB name as TABLE_CAT in JDBC
metadata results, but the TableId has a null catalog. Fixed by using a
catalog-agnostic filter when calling readSchema().
6. {*}{{*}}NullPointerException on absent ROWID Connect header{{*}}{*}:
Debezium 3.4.2 removed
the getEmitConnectHeaders() override that injected a ROWID into Kafka
Connect
headers. OracleSourceFetchTaskContext.isRecordBetween() then crashes on
tables
without a primary key because headers.iterator().next() throws
NoSuchElementException when no header is present.
7. {*}{{*}}ClassNotFoundException at runtime{{*}}{*}: debezium-common classes
(e.g.
RelationalDatabaseSchema) are referenced at runtime but were not included
in
the Maven shade assembly, causing ClassNotFoundException when the fat JAR
is used.
*Steps to Reproduce*
1. Set up Oracle 26ai in CDB/PDB mode (e.g. FREE CDB + FREEPDB1)
2. Upgrade Debezium version in pom.xml to 3.4.2.Final
3. Run Oracle CDC connector with scan.startup.mode=initial
4. Observe failures listed above during snapshot and streaming phases
*Fix*
- Update vendored LogMinerChangeRecordEmitter to 9-arg constructor matching
Debezium 3.4.2
- Add local ResumePositionProvider override: null-guard on sessionContext in
close() and
catch block; strip "database." prefix from JdbcConfiguration before
creating OracleConnection
- Fix OracleUtils.queryApproximateRowCnt() to use ALL_TABLES.NUM_ROWS instead
of ANALYZE TABLE
- Fix OracleSchema.getTableSchema() to use a catalog-agnostic TableFilter for
readSchema()
so it works correctly in Oracle CDB/PDB mode
- Guard isRecordBetween() against absent ROWID Connect header; fall back to
SCN-range
inclusion for tables without a primary key (multi-chunk PK-less tables
should set
scan.incremental.snapshot.chunk.key-column explicitly)
- Include debezium-common in Maven shade assembly to prevent
ClassNotFoundException at runtime
- Add EmbeddedFlinkSchemaHistory and RelationalDatabaseSchema adapters for
Debezium 3.x API
*Verification*
Tested end-to-end pipeline: Oracle 26ai FREE (CDB/FREEPDB1) → Flink CDC
3.6-SNAPSHOT
→ StarRocks sink. Confirmed 3 snapshot rows and a live CDC UPDATE propagated
correctly
using Debezium 3.4.2.Final and scan.startup.mode=initial.
*Environment*
- Flink CDC: 3.6-SNAPSHOT
- Debezium: 3.4.2.Final
- Oracle: 26ai (FREE edition, CDB + FREEPDB1 PDB)
- Java: 17
> [Oracle CDC] Fix compatibility with Oracle 26ai and Debezium 3.4.2.Final
> ------------------------------------------------------------------------
>
> Key: FLINK-40201
> URL: https://issues.apache.org/jira/browse/FLINK-40201
> Project: Flink
> Issue Type: Bug
> Components: Flink CDC
> Affects Versions: cdc-3.6.0
> Reporter: Mohammad Suhel Ksolves Cassandra Experties
> Priority: Major
>
> Here it is — copy and paste directly:
> ----
> h2. Background
> Flink CDC 3.6's Oracle connector vendors several Debezium internal classes
> that were written against an older Debezium API. Oracle 26ai (23.26.x)
> requires features only available in Debezium 3.4.2.Final (specifically
> UnbufferedLogMinerStreamingChangeEventSource introduced in Debezium 2.x).
> When Debezium 3.4.2.Final is used, the vendored classes produce runtime
> errors because they still carry old API signatures. Simply swapping the
> Debezium JAR is not sufficient — several source-level fixes are required to
> align the vendored classes with the Debezium 3.4.2 API.
> h2. Problem
> Oracle CDC connector in Flink CDC 3.6 fails with Oracle 26ai (CDB/PDB mode)
> when using Debezium 3.4.2.Final. Multiple runtime errors prevent snapshot and
> streaming CDC from working end-to-end:
> # {*}NoSuchMethodError on LogMinerChangeRecordEmitter{*}: The vendored
> LogMinerChangeRecordEmitter carries a 10-arg constructor (from an older
> Debezium API) but Debezium 3.4.2's
> UnbufferedLogMinerStreamingChangeEventSource calls a 9-arg constructor (rowId
> parameter was removed in Debezium 2.x).
> # {*}NullPointerException in ResumePositionProvider.close(){*}: The close()
> method calls sessionContext.close() without null-guarding sessionContext. If
> lazy initialisation fails after connection is assigned but before
> sessionContext is assigned, close() produces a secondary NPE that hides the
> root cause exception.
> # {*}Invalid JDBC connection string in ResumePositionProvider{*}: When the
> connector config carries "database.*" prefixed keys, OracleConnection
> receives keys like "database.hostname" instead of "hostname", producing a
> malformed JDBC URL and connection failure.
> # {*}ORA-01031 on row count query{*}: OracleUtils.queryApproximateRowCnt()
> used ANALYZE TABLE which requires the ANALYZE ANY system privilege. Standard
> CDC users do not have this privilege. Fixed by querying ALL_TABLES.NUM_ROWS
> instead.
> # {*}ORA-00942 / schema recovery failure{*}: readSchemaForCapturedTables()
> fails in CDB/PDB mode because Oracle JDBC returns the PDB name as TABLE_CAT
> in JDBC metadata results, but the TableId has a null catalog. Fixed by using
> a catalog-agnostic filter when calling readSchema().
> # {*}NullPointerException on absent ROWID Connect header{*}: Debezium 3.4.2
> removed the getEmitConnectHeaders() override that injected a ROWID into Kafka
> Connect headers. OracleSourceFetchTaskContext.isRecordBetween() then crashes
> on tables without a primary key because headers.iterator().next() throws
> NoSuchElementException when no header is present.
> # {*}ClassNotFoundException at runtime{*}: debezium-common classes (e.g.
> RelationalDatabaseSchema) are referenced at runtime but were not included in
> the Maven shade assembly, causing ClassNotFoundException when the fat JAR is
> used.
> h2. Steps to Reproduce
> # Set up Oracle 26ai in CDB/PDB mode (e.g. FREE CDB + FREEPDB1)
> # Upgrade Debezium version in pom.xml to 3.4.2.Final
> # Run Oracle CDC connector with scan.startup.mode=initial
> # Observe failures listed above during snapshot and streaming phases
> h2. Fix
> * Update vendored LogMinerChangeRecordEmitter to 9-arg constructor matching
> Debezium 3.4.2
> * Add local ResumePositionProvider override: null-guard on sessionContext in
> close() and catch block; strip "database." prefix from JdbcConfiguration
> before creating OracleConnection
> * Fix OracleUtils.queryApproximateRowCnt() to use ALL_TABLES.NUM_ROWS
> instead of ANALYZE TABLE
> * Fix OracleSchema.getTableSchema() to use a catalog-agnostic TableFilter
> for readSchema() so it works correctly in Oracle CDB/PDB mode
> * Guard isRecordBetween() against absent ROWID Connect header; fall back to
> SCN-range inclusion for tables without a primary key (multi-chunk PK-less
> tables should set scan.incremental.snapshot.chunk.key-column explicitly)
> * Include debezium-common in Maven shade assembly to prevent
> ClassNotFoundException at runtime
> * Add EmbeddedFlinkSchemaHistory and RelationalDatabaseSchema adapters for
> Debezium 3.x API
> h2. Verification
> Tested end-to-end pipeline: Oracle 26ai FREE (CDB/FREEPDB1) → Flink CDC
> 3.6-SNAPSHOT → StarRocks sink. Confirmed 3 snapshot rows and a live CDC
> UPDATE propagated correctly using Debezium 3.4.2.Final and
> scan.startup.mode=initial.
> h2. Environment
> * Flink CDC: 3.6-SNAPSHOT
> * Debezium: 3.4.2.Final
> * Oracle: 26ai (FREE edition, CDB + FREEPDB1 PDB)
> * Java: 17
--
This message was sent by Atlassian Jira
(v8.20.10#820010)