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 708a7b0645 [Fix][Connector-V2][JDBC] Apply where_condition to JDBC 
split metadata queries (#12160)
708a7b0645 is described below

commit 708a7b06457ab4f3ece70e45e9a21334abfc9c8d
Author: Jast <[email protected]>
AuthorDate: Fri Sep 11 02:49:54 2026 +0000

    [Fix][Connector-V2][JDBC] Apply where_condition to JDBC split metadata 
queries (#12160)
    
    Co-authored-by: zhangshenghang <[email protected]>
---
 .../seatunnel/jdbc/source/ChunkSplitter.java       | 44 ++++++++++-
 .../jdbc/source/DynamicChunkSplitter.java          |  7 +-
 .../seatunnel/jdbc/source/JdbcSourceReader.java    | 13 ++++
 .../seatunnel/jdbc/source/JdbcSourceTable.java     | 22 ++++++
 .../jdbc/source/DynamicChunkSplitterTest.java      | 85 ++++++++++++++++++++++
 5 files changed, 165 insertions(+), 6 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/ChunkSplitter.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/ChunkSplitter.java
index 2d6282e2f8..d2499ea249 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/ChunkSplitter.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/ChunkSplitter.java
@@ -52,6 +52,8 @@ import java.util.stream.Collectors;
 @Slf4j
 public abstract class ChunkSplitter implements AutoCloseable, Serializable {
 
+    private static final int SPLIT_COUNT_WARN_THRESHOLD = 512;
+
     protected JdbcSourceConfig config;
     protected final JdbcConnectionProvider connectionProvider;
     protected final JdbcDialect jdbcDialect;
@@ -124,10 +126,20 @@ public abstract class ChunkSplitter implements 
AutoCloseable, Serializable {
 
         long end = System.currentTimeMillis();
         log.info(
-                "Split table {} into {} chunks, time cost: {}ms.",
+                "Split table {} into {} chunks, time cost: {}ms. Where 
condition configured for split metadata queries: {}",
                 table.getTablePath(),
                 splits.size(),
-                end - start);
+                end - start,
+                StringUtils.isNotBlank(config.getWhereConditionClause()));
+        if (splits.size() > SPLIT_COUNT_WARN_THRESHOLD) {
+            log.warn(
+                    "Table {} was split into {} chunks, above the warning 
threshold {}. "
+                            + "A large number of mostly-empty chunks slows the 
job down; "
+                            + "check the where condition, split key and chunk 
size.",
+                    table.getTablePath(),
+                    splits.size(),
+                    SPLIT_COUNT_WARN_THRESHOLD);
+        }
         return splits;
     }
 
@@ -206,7 +218,7 @@ public abstract class ChunkSplitter implements 
AutoCloseable, Serializable {
 
         StringRangeSplitDecision decision =
                 jdbcDialect.validateStringRangeSplit(
-                        getOrEstablishConnection(), table, splitKeyName, 256);
+                        getOrEstablishConnection(), 
applyWhereCondition(table), splitKeyName, 256);
         if (decision.isSafe()) {
             return StringSplitStrategy.RANGE;
         }
@@ -244,8 +256,33 @@ public abstract class ChunkSplitter implements 
AutoCloseable, Serializable {
         return createPreparedStatement(splitQuery);
     }
 
+    /**
+     * Wraps the table query with the configured where condition so that split 
metadata queries
+     * (min/max, row count, chunk boundary, sampling) run on the same data 
scope as the split reads,
+     * which apply the where condition separately in {@link 
#createPreparedStatement}. Reuses {@link
+     * SqlWhereConditionHelper#applyWhereConditionWithWrap} so 
where-referenced columns missing from
+     * a narrow custom query projection are auto-added, exactly like the read 
path. Returns the
+     * table unchanged when no where condition is configured.
+     */
+    protected JdbcSourceTable applyWhereCondition(JdbcSourceTable table) {
+        if (StringUtils.isBlank(config.getWhereConditionClause())) {
+            return table;
+        }
+        String baseQuery =
+                StringUtils.isNotBlank(table.getQuery())
+                        ? table.getQuery()
+                        : String.format(
+                                "SELECT * FROM %s",
+                                
jdbcDialect.tableIdentifier(table.getTablePath()));
+        String effectiveQuery =
+                SqlWhereConditionHelper.applyWhereConditionWithWrap(
+                        baseQuery, config.getWhereConditionClause(), true);
+        return table.withQuery(effectiveQuery);
+    }
+
     protected Object queryMin(JdbcSourceTable table, String columnName, Object 
excludedLowerBound)
             throws SQLException {
+        table = applyWhereCondition(table);
         String minQuery;
         Map<String, Column> columns =
                 table.getCatalogTable().getTableSchema().getColumns().stream()
@@ -285,6 +322,7 @@ public abstract class ChunkSplitter implements 
AutoCloseable, Serializable {
 
     protected Pair<Object, Object> queryMinMax(JdbcSourceTable table, String 
columnName)
             throws SQLException {
+        table = applyWhereCondition(table);
         String sqlQuery;
         Map<String, Column> columns =
                 table.getCatalogTable().getTableSchema().getColumns().stream()
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitter.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitter.java
index 12c50fc66f..7660d77205 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitter.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitter.java
@@ -497,7 +497,7 @@ public class DynamicChunkSplitter extends ChunkSplitter {
             Object[] sample =
                     jdbcDialect.sampleDataFromColumn(
                             getOrEstablishConnection(),
-                            table,
+                            applyWhereCondition(table),
                             splitColumnName,
                             inverseSamplingRate,
                             config.getFetchSize());
@@ -512,7 +512,8 @@ public class DynamicChunkSplitter extends ChunkSplitter {
     }
 
     private Long queryApproximateRowCnt(JdbcSourceTable table) throws 
SQLException {
-        return 
jdbcDialect.approximateRowCntStatement(getOrEstablishConnection(), table);
+        return jdbcDialect.approximateRowCntStatement(
+                getOrEstablishConnection(), applyWhereCondition(table));
     }
 
     private double calculateDistributionFactor(
@@ -891,7 +892,7 @@ public class DynamicChunkSplitter extends ChunkSplitter {
         Object chunkEnd =
                 jdbcDialect.queryNextChunkMax(
                         getOrEstablishConnection(),
-                        table,
+                        applyWhereCondition(table),
                         splitColumnName,
                         chunkSize,
                         previousChunkEnd);
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceReader.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceReader.java
index f1b85dbfc7..40bbe28d97 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceReader.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceReader.java
@@ -33,13 +33,18 @@ import java.util.Deque;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.ConcurrentLinkedDeque;
+import java.util.concurrent.atomic.AtomicInteger;
 
 @Slf4j
 public class JdbcSourceReader implements SourceReader<SeaTunnelRow, 
JdbcSourceSplit> {
+    private static final int SPLIT_PROGRESS_LOG_INTERVAL = 50;
+
     private final Context context;
     private final JdbcInputFormat inputFormat;
     private final Deque<JdbcSourceSplit> splits = new 
ConcurrentLinkedDeque<>();
     private volatile boolean noMoreSplit;
+    private final AtomicInteger assignedSplitCount = new AtomicInteger();
+    private final AtomicInteger processedSplitCount = new AtomicInteger();
 
     public JdbcSourceReader(
             Context context, JdbcSourceConfig config, Map<TablePath, 
CatalogTable> tables) {
@@ -72,6 +77,13 @@ public class JdbcSourceReader implements 
SourceReader<SeaTunnelRow, JdbcSourceSp
                 } finally {
                     inputFormat.close();
                 }
+                int processedCount = processedSplitCount.incrementAndGet();
+                if (processedCount % SPLIT_PROGRESS_LOG_INTERVAL == 0) {
+                    log.info(
+                            "Processed {} of {} assigned jdbc source splits",
+                            processedCount,
+                            assignedSplitCount.get());
+                }
             } else if (noMoreSplit && splits.isEmpty()) {
                 // signal to the source that we have reached the end of the 
data.
                 log.info("Closed the bounded jdbc source");
@@ -90,6 +102,7 @@ public class JdbcSourceReader implements 
SourceReader<SeaTunnelRow, JdbcSourceSp
     @Override
     public void addSplits(List<JdbcSourceSplit> splits) {
         this.splits.addAll(splits);
+        this.assignedSplitCount.addAndGet(splits.size());
     }
 
     @Override
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceTable.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceTable.java
index 7144c9ed2d..2db70895eb 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceTable.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceTable.java
@@ -39,4 +39,26 @@ public class JdbcSourceTable implements Serializable {
     private final Boolean useSelectCount;
     private final Boolean skipAnalyze;
     private final CatalogTable catalogTable;
+
+    /**
+     * Returns a copy of this table with its query replaced by {@code query}, 
preserving every other
+     * field. Used to scope split-metadata queries to a 
where-condition-filtered view of the table
+     * without mutating the original.
+     *
+     * <p>NOTE: keep this manual copy-constructor in sync whenever a new field 
is added to this
+     * class, otherwise the new field will be silently lost in the copy.
+     */
+    public JdbcSourceTable withQuery(String query) {
+        return JdbcSourceTable.builder()
+                .tablePath(this.tablePath)
+                .query(query)
+                .partitionColumn(this.partitionColumn)
+                .partitionNumber(this.partitionNumber)
+                .partitionStart(this.partitionStart)
+                .partitionEnd(this.partitionEnd)
+                .useSelectCount(this.useSelectCount)
+                .skipAnalyze(this.skipAnalyze)
+                .catalogTable(this.catalogTable)
+                .build();
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitterTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitterTest.java
index c9fe3c2f1c..dba28f82b6 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitterTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitterTest.java
@@ -38,6 +38,7 @@ import java.util.Map;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
 public class DynamicChunkSplitterTest {
@@ -387,4 +388,88 @@ public class DynamicChunkSplitterTest {
             }
         }
     }
+
+    /** Without a where condition the table must pass through unchanged. */
+    @Test
+    public void testApplyWhereConditionReturnsSameTableWhenNoWhereCondition() {
+        JdbcSourceConfig config = buildWhereConditionConfig(null);
+        DynamicChunkSplitter splitter = new DynamicChunkSplitter(config);
+        JdbcSourceTable table =
+                JdbcSourceTable.builder()
+                        .tablePath(TablePath.of("db", "schema", "table"))
+                        .query("SELECT id, name FROM table")
+                        .build();
+
+        assertSame(table, splitter.applyWhereCondition(table));
+    }
+
+    /** The user query must be wrapped with the where condition for split 
metadata queries. */
+    @Test
+    public void testApplyWhereConditionWrapsUserQuery() {
+        JdbcSourceConfig config = buildWhereConditionConfig("where id > 100");
+        DynamicChunkSplitter splitter = new DynamicChunkSplitter(config);
+        JdbcSourceTable table =
+                JdbcSourceTable.builder()
+                        .tablePath(TablePath.of("db", "schema", "table"))
+                        .query("SELECT id, name FROM table")
+                        .build();
+
+        JdbcSourceTable wrapped = splitter.applyWhereCondition(table);
+
+        assertEquals(
+                "SELECT * FROM (SELECT id, name FROM table) tmp WHERE id > 
100",
+                wrapped.getQuery());
+        assertSame(table.getTablePath(), wrapped.getTablePath());
+        assertEquals(table.getPartitionColumn(), wrapped.getPartitionColumn());
+        assertSame(table.getCatalogTable(), wrapped.getCatalogTable());
+    }
+
+    /** A where-referenced column missing from a narrow custom query must be 
auto-added. */
+    @Test
+    public void testApplyWhereConditionAutoAddsMissingFieldForNarrowQuery() {
+        JdbcSourceConfig config = buildWhereConditionConfig("where status > 
1");
+        DynamicChunkSplitter splitter = new DynamicChunkSplitter(config);
+        JdbcSourceTable table =
+                JdbcSourceTable.builder()
+                        .tablePath(TablePath.of("db", "schema", "table"))
+                        .query("SELECT id, name FROM table")
+                        .build();
+
+        JdbcSourceTable wrapped = splitter.applyWhereCondition(table);
+
+        // Without the auto-add, the wrapped subquery would not expose 
"status" and the
+        // split-metadata queries would fail with a "column not found" SQL 
error.
+        assertEquals(
+                "SELECT * FROM (SELECT id, name , status FROM table) tmp WHERE 
status > 1",
+                wrapped.getQuery());
+    }
+
+    /** When no query is configured the table identifier must be used as the 
wrapped base. */
+    @Test
+    public void 
testApplyWhereConditionFallsBackToTableIdentifierWithoutQuery() {
+        JdbcSourceConfig config = buildWhereConditionConfig("where id > 100");
+        DynamicChunkSplitter splitter = new DynamicChunkSplitter(config);
+        JdbcSourceTable table =
+                JdbcSourceTable.builder().tablePath(TablePath.of("db", 
"schema", "table")).build();
+
+        JdbcSourceTable wrapped = splitter.applyWhereCondition(table);
+
+        // The base query is "SELECT * FROM <tableIdentifier>"; for the 
default Postgres
+        // dialect the tableIdentifier is the fully quoted path. We only 
assert that the
+        // wrapper is applied and the original table path is preserved.
+        assertEquals(
+                "SELECT * FROM (SELECT * FROM \"db\".\"schema\".\"table\") tmp 
WHERE id > 100",
+                wrapped.getQuery());
+        assertSame(table.getTablePath(), wrapped.getTablePath());
+    }
+
+    private static JdbcSourceConfig buildWhereConditionConfig(String 
whereCondition) {
+        Map<String, Object> options = new HashMap<>();
+        options.put("url", "jdbc:postgresql://localhost:5432/test");
+        options.put("driver", "org.postgresql.Driver");
+        if (whereCondition != null) {
+            options.put(JdbcSourceOptions.WHERE_CONDITION.key(), 
whereCondition);
+        }
+        return JdbcSourceConfig.of(ReadonlyConfig.fromMap(options));
+    }
 }

Reply via email to