This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 70faa7ce50 [Fix][Connector-V2] Honor Oracle snapshot select overrides 
(#11768)
70faa7ce50 is described below

commit 70faa7ce50c62793fc2ade19df288b5219774af6
Author: Jast <[email protected]>
AuthorDate: Fri Sep 11 02:49:55 2026 +0000

    [Fix][Connector-V2] Honor Oracle snapshot select overrides (#11768)
    
    Co-authored-by: zhangshenghang <[email protected]>
---
 docs/en/connectors/source/Oracle-CDC.md            | 11 +++
 docs/zh/connectors/source/Oracle-CDC.md            | 11 +++
 .../fetch/scan/OracleSnapshotSplitReadTask.java    |  3 +-
 .../seatunnel/cdc/oracle/utils/OracleUtils.java    | 82 ++++++++++++++++------
 .../cdc/oracle/utils/OracleUtilsTest.java          | 29 ++++++++
 5 files changed, 112 insertions(+), 24 deletions(-)

diff --git a/docs/en/connectors/source/Oracle-CDC.md 
b/docs/en/connectors/source/Oracle-CDC.md
index df9bcb9b0e..941f247a95 100644
--- a/docs/en/connectors/source/Oracle-CDC.md
+++ b/docs/en/connectors/source/Oracle-CDC.md
@@ -564,6 +564,17 @@ Yes. Set `database-names` to the CDB name and configure 
the JDBC URL to point to
 
 By default, Oracle CDC requires primary keys. You can specify a custom primary 
key column via `table-names-config` with the `primaryKeys` field if the table 
has a suitable unique column.
 
+### How do I use a custom snapshot query?
+
+Use Debezium's `snapshot.select.statement.overrides` properties inside the 
`debezium` block. The query is applied before SeaTunnel adds the snapshot-split 
boundaries, so it must include every column needed by the configured table 
schema and its split key.
+
+```hocon
+debezium {
+  snapshot.select.statement.overrides = "DEBEZIUM.FULL_TYPES"
+  snapshot.select.statement.overrides.DEBEZIUM.FULL_TYPES = "SELECT * FROM 
DEBEZIUM.FULL_TYPES WHERE ACTIVE = 1"
+}
+```
+
 ### How do I improve LogMiner performance?
 
 Treat this primarily as a database and redo-log tuning topic. Reuse the 
LogMiner setup and
diff --git a/docs/zh/connectors/source/Oracle-CDC.md 
b/docs/zh/connectors/source/Oracle-CDC.md
index 41c52fe2d4..1313323df2 100644
--- a/docs/zh/connectors/source/Oracle-CDC.md
+++ b/docs/zh/connectors/source/Oracle-CDC.md
@@ -558,6 +558,17 @@ ALTER TABLE schema_name.table_name ADD SUPPLEMENTAL LOG 
DATA (ALL) COLUMNS;
 
 默认情况下,Oracle CDC 需要主键。如果表中存在合适的唯一列,可通过 `table-names-config` 中的 `primaryKeys` 
字段指定自定义主键列。
 
+### 如何使用自定义快照查询?
+
+在 `debezium` 块中配置 Debezium 的 `snapshot.select.statement.overrides` 
属性。SeaTunnel 会先使用该查询,再追加快照分片边界条件,因此查询必须包含已配置表结构和分片键所需的全部列。
+
+```hocon
+debezium {
+  snapshot.select.statement.overrides = "DEBEZIUM.FULL_TYPES"
+  snapshot.select.statement.overrides.DEBEZIUM.FULL_TYPES = "SELECT * FROM 
DEBEZIUM.FULL_TYPES WHERE ACTIVE = 1"
+}
+```
+
 ### 如何提升 LogMiner 性能?
 
 首先把它当作数据库和 redo log 调优问题处理。优先复用上面的 LogMiner 配置和 supplemental
diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/reader/fetch/scan/OracleSnapshotSplitReadTask.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/reader/fetch/scan/OracleSnapshotSplitReadTask.java
index b90b0361fb..98a9977728 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/reader/fetch/scan/OracleSnapshotSplitReadTask.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/reader/fetch/scan/OracleSnapshotSplitReadTask.java
@@ -206,7 +206,8 @@ public class OracleSnapshotSplitReadTask
             LOG.info("Switched to PDB '{}' for table {}", 
connectorConfig.getPdbName(), table.id());
         }
         final String selectSql =
-                OracleUtils.buildSplitScanQuery(
+                OracleUtils.buildSnapshotSplitScanQuery(
+                        connectorConfig,
                         snapshotSplit.getTableId(),
                         snapshotSplit.getSplitKeyType(),
                         snapshotSplit.getSplitStart() == null,
diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtils.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtils.java
index fbb3664be0..4c161392c7 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtils.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtils.java
@@ -255,6 +255,37 @@ public class OracleUtils {
         return buildSplitQuery(tableId, rowType, isFirstSplit, isLastSplit, 
-1, true);
     }
 
+    /**
+     * Builds the snapshot query for a split, honoring a Debezium per-table 
select override when
+     * configured.
+     *
+     * <p>The override is wrapped as an inline view before the split predicate 
is appended. This
+     * retains the caller's filter while keeping SeaTunnel's parallel snapshot 
boundaries intact.
+     */
+    public static String buildSnapshotSplitScanQuery(
+            OracleConnectorConfig connectorConfig,
+            TableId tableId,
+            SeaTunnelRowType rowType,
+            boolean isFirstSplit,
+            boolean isLastSplit) {
+        String overriddenSelect = 
connectorConfig.getSnapshotSelectOverridesByTable().get(tableId);
+        if (overriddenSelect == null) {
+            overriddenSelect =
+                    connectorConfig
+                            .getSnapshotSelectOverridesByTable()
+                            .get(new TableId(null, tableId.schema(), 
tableId.table()));
+        }
+        if (overriddenSelect == null) {
+            return buildSplitScanQuery(tableId, rowType, isFirstSplit, 
isLastSplit);
+        }
+
+        String condition = buildSplitCondition(rowType, isFirstSplit, 
isLastSplit, true);
+        if (condition == null) {
+            return overriddenSelect;
+        }
+        return String.format("SELECT * FROM (%s) WHERE %s", overriddenSelect, 
condition);
+    }
+
     private static String buildSplitQuery(
             TableId tableId,
             SeaTunnelRowType rowType,
@@ -262,25 +293,44 @@ public class OracleUtils {
             boolean isLastSplit,
             int limitSize,
             boolean isScanningData) {
-        final String condition;
+        final String condition =
+                buildSplitCondition(rowType, isFirstSplit, isLastSplit, 
isScanningData);
 
+        if (isScanningData) {
+            return buildSelectWithRowLimits(
+                    tableId, limitSize, "*", Optional.ofNullable(condition), 
Optional.empty());
+        } else {
+            final String orderBy = String.join(", ", rowType.getFieldNames());
+            return buildSelectWithBoundaryRowLimits(
+                    tableId,
+                    limitSize,
+                    getPrimaryKeyColumnsProjection(rowType),
+                    getMaxPrimaryKeyColumnsProjection(rowType),
+                    Optional.ofNullable(condition),
+                    orderBy);
+        }
+    }
+
+    private static String buildSplitCondition(
+            SeaTunnelRowType rowType,
+            boolean isFirstSplit,
+            boolean isLastSplit,
+            boolean isScanningData) {
         if (isFirstSplit && isLastSplit) {
-            condition = null;
-        } else if (isFirstSplit) {
-            final StringBuilder sql = new StringBuilder();
+            return null;
+        }
+
+        final StringBuilder sql = new StringBuilder();
+        if (isFirstSplit) {
             addPrimaryKeyColumnsToCondition(rowType, sql, " <= ?");
             if (isScanningData) {
                 sql.append(" AND NOT (");
                 addPrimaryKeyColumnsToCondition(rowType, sql, " = ?");
                 sql.append(")");
             }
-            condition = sql.toString();
         } else if (isLastSplit) {
-            final StringBuilder sql = new StringBuilder();
             addPrimaryKeyColumnsToCondition(rowType, sql, " >= ?");
-            condition = sql.toString();
         } else {
-            final StringBuilder sql = new StringBuilder();
             addPrimaryKeyColumnsToCondition(rowType, sql, " >= ?");
             if (isScanningData) {
                 sql.append(" AND NOT (");
@@ -289,22 +339,8 @@ public class OracleUtils {
             }
             sql.append(" AND ");
             addPrimaryKeyColumnsToCondition(rowType, sql, " <= ?");
-            condition = sql.toString();
-        }
-
-        if (isScanningData) {
-            return buildSelectWithRowLimits(
-                    tableId, limitSize, "*", Optional.ofNullable(condition), 
Optional.empty());
-        } else {
-            final String orderBy = String.join(", ", rowType.getFieldNames());
-            return buildSelectWithBoundaryRowLimits(
-                    tableId,
-                    limitSize,
-                    getPrimaryKeyColumnsProjection(rowType),
-                    getMaxPrimaryKeyColumnsProjection(rowType),
-                    Optional.ofNullable(condition),
-                    orderBy);
         }
+        return sql.toString();
     }
 
     public static PreparedStatement readTableSplitDataStatement(
diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtilsTest.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtilsTest.java
index a3856a7ff7..2173289556 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtilsTest.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/utils/OracleUtilsTest.java
@@ -24,6 +24,8 @@ import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import io.debezium.config.Configuration;
+import io.debezium.connector.oracle.OracleConnectorConfig;
 import io.debezium.relational.TableId;
 
 import java.util.Collections;
@@ -74,6 +76,33 @@ public class OracleUtilsTest {
                 "SELECT * FROM \"schema1\".\"table1\" WHERE \"id\" >= ?", 
splitScanSQL);
     }
 
+    @Test
+    public void testSnapshotSplitScanQueryUsesSelectOverride() {
+        TableId tableId = TableId.parse("cdb1.schema1.table1");
+        SeaTunnelRowType splitKeyType =
+                new SeaTunnelRowType(
+                        new String[] {"id"}, new SeaTunnelDataType[] 
{BasicType.LONG_TYPE});
+        String overriddenSelect = "SELECT id, name FROM schema1.table1 WHERE 
active = 1";
+        OracleConnectorConfig connectorConfig =
+                new OracleConnectorConfig(
+                        Configuration.create()
+                                .with(OracleConnectorConfig.SERVER_NAME, 
"test_server")
+                                .with(OracleConnectorConfig.HOSTNAME, 
"localhost")
+                                .with(OracleConnectorConfig.USER, "test")
+                                .with(OracleConnectorConfig.PASSWORD, "test")
+                                .with("snapshot.select.statement.overrides", 
"schema1.table1")
+                                .with(
+                                        
"snapshot.select.statement.overrides.schema1.table1",
+                                        overriddenSelect)
+                                .build());
+
+        Assertions.assertEquals(
+                "SELECT * FROM (SELECT id, name FROM schema1.table1 WHERE 
active = 1) "
+                        + "WHERE \"id\" >= ? AND NOT (\"id\" = ?) AND \"id\" 
<= ?",
+                OracleUtils.buildSnapshotSplitScanQuery(
+                        connectorConfig, tableId, splitKeyType, false, false));
+    }
+
     @Test
     public void testResolveTableIdWithRequestedCatalog() {
         TableId requestedTableId = TableId.parse("ORCLPDB.LIB_B.T_B1");

Reply via email to