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");