From 4a6337e6b7f4cb9039027000d446f18fd9251c67 Mon Sep 17 00:00:00 2001 From: rockyyin Date: Tue, 18 Aug 2026 15:48:04 +0800 Subject: [PATCH] [oracle] Support specific-offset startup mode with SCN for Oracle pipeline source The Oracle pipeline connector only supports initial/snapshot/latest-offset startup modes, while the MySQL connector already supports specific-offset. This adds specific-offset mode to the Oracle DataSource, allowing users to start incremental reading (LogMiner redo log consumption) from a specified SCN via scan.startup.specific-offset.scn. Changes: add SCAN_STARTUP_SPECIFIC_OFFSET_SCN option in OracleDataSourceOptions; add specific-offset branch in getStartupOptions() of OracleDataSourceFactory, building StartupOptions.specificOffset with {scn, commit_scn} map which is natively consumed by the base framework StreamSplitAssigner. --- .../factory/OracleDataSourceFactory.java | 22 ++++++++++++++++++- .../source/OracleDataSourceOptions.java | 10 ++++++++- 2 files changed, 30 insertions(+), 2 deletions(-) diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/factory/OracleDataSourceFactory.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/factory/OracleDataSourceFactory.java index a2b4647fc72..aabb3d20817 100644 --- a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/factory/OracleDataSourceFactory.java +++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/factory/OracleDataSourceFactory.java @@ -33,6 +33,7 @@ import org.apache.flink.cdc.connectors.oracle.source.OracleDataSourceOptions; import org.apache.flink.cdc.connectors.oracle.source.config.OracleSourceConfig; import org.apache.flink.cdc.connectors.oracle.source.config.OracleSourceConfigFactory; +import org.apache.flink.cdc.connectors.oracle.source.meta.offset.RedoLogOffset; import org.apache.flink.cdc.connectors.oracle.table.OracleReadableMetaData; import org.apache.flink.cdc.connectors.oracle.utils.OracleSchemaUtils; import org.apache.flink.table.api.ValidationException; @@ -40,6 +41,7 @@ import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; +import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -50,6 +52,7 @@ import static org.apache.flink.cdc.connectors.base.utils.ObjectUtils.doubleCompare; import static org.apache.flink.cdc.connectors.oracle.source.OracleDataSourceOptions.METADATA_LIST; import static org.apache.flink.cdc.connectors.oracle.source.OracleDataSourceOptions.SCAN_STARTUP_MODE; +import static org.apache.flink.cdc.connectors.oracle.source.OracleDataSourceOptions.SCAN_STARTUP_SPECIFIC_OFFSET_SCN; import static org.apache.flink.cdc.debezium.table.DebeziumOptions.DEBEZIUM_OPTIONS_PREFIX; import static org.apache.flink.cdc.debezium.utils.JdbcUrlUtils.PROPERTIES_PREFIX; import static org.apache.flink.util.Preconditions.checkNotNull; @@ -63,6 +66,7 @@ public class OracleDataSourceFactory implements DataSourceFactory { private static final String SCAN_STARTUP_MODE_VALUE_INITIAL = "initial"; private static final String SCAN_STARTUP_MODE_VALUE_LATEST = "latest-offset"; private static final String SCAN_STARTUP_MODE_VALUE_SNAPSHOT = "snapshot"; + private static final String SCAN_STARTUP_MODE_VALUE_SPECIFIC = "specific-offset"; @Override public DataSource createDataSource(Context context) { @@ -204,6 +208,7 @@ public Set> optionalOptions() { options.add(OracleDataSourceOptions.LOG_MINING_STRATEGY); options.add(OracleDataSourceOptions.DATABASE_CONNECTION_ADAPTER); options.add(OracleDataSourceOptions.SCAN_STARTUP_MODE); + options.add(OracleDataSourceOptions.SCAN_STARTUP_SPECIFIC_OFFSET_SCN); return options; } @@ -265,15 +270,30 @@ private static StartupOptions getStartupOptions(Configuration config) { return StartupOptions.snapshot(); case SCAN_STARTUP_MODE_VALUE_LATEST: return StartupOptions.latest(); + case SCAN_STARTUP_MODE_VALUE_SPECIFIC: + Long scn = config.get(SCAN_STARTUP_SPECIFIC_OFFSET_SCN); + if (scn == null) { + throw new ValidationException( + String.format( + "Option '%s' is required when '%s' is '%s'.", + SCAN_STARTUP_SPECIFIC_OFFSET_SCN.key(), + SCAN_STARTUP_MODE.key(), + SCAN_STARTUP_MODE_VALUE_SPECIFIC)); + } + Map specificOffsetMap = new HashMap<>(); + specificOffsetMap.put(RedoLogOffset.SCN_KEY, String.valueOf(scn)); + specificOffsetMap.put(RedoLogOffset.COMMIT_SCN_KEY, String.valueOf(0L)); + return StartupOptions.specificOffset(specificOffsetMap); default: throw new ValidationException( String.format( - "Invalid value for option '%s'. Supported values are [%s, %s, %s], but was: %s", + "Invalid value for option '%s'. Supported values are [%s, %s, %s, %s], but was: %s", SourceOptions.SCAN_STARTUP_MODE.key(), SCAN_STARTUP_MODE_VALUE_INITIAL, SCAN_STARTUP_MODE_VALUE_SNAPSHOT, SCAN_STARTUP_MODE_VALUE_LATEST, + SCAN_STARTUP_MODE_VALUE_SPECIFIC, modeString)); } } diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/source/OracleDataSourceOptions.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/source/OracleDataSourceOptions.java index b0f1296250e..117829881b5 100644 --- a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/source/OracleDataSourceOptions.java +++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oracle/src/main/java/org/apache/flink/cdc/connectors/oracle/source/OracleDataSourceOptions.java @@ -130,7 +130,15 @@ public class OracleDataSourceOptions { .defaultValue("initial") .withDescription( "Optional startup mode for oracle CDC consumer, valid enumerations are " - + "\"initial\", \"latest-offset\", \"snapshot\""); + + "\"initial\", \"latest-offset\", \"snapshot\", \"specific-offset\""); + + public static final ConfigOption SCAN_STARTUP_SPECIFIC_OFFSET_SCN = + ConfigOptions.key("scan.startup.specific-offset.scn") + .longType() + .noDefaultValue() + .withDescription( + "The specific SCN to start reading redo log from, " + + "required when 'scan.startup.mode' is 'specific-offset'."); @Experimental public static final ConfigOption CHUNK_META_GROUP_SIZE =