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));
+ }
}