This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new eab7728694 [core] Support DV-aware TopN pushdown (#8363)
eab7728694 is described below
commit eab7728694ad0706c20006865a7f0b3b60db3814
Author: Kerwin Zhang <[email protected]>
AuthorDate: Sat Jun 27 22:41:08 2026 +0800
[core] Support DV-aware TopN pushdown (#8363)
---
.../paimon/table/source/DataTableBatchScan.java | 5 +-
.../table/source/TopNDataSplitEvaluator.java | 4 +-
.../apache/paimon/table/source/TableScanTest.java | 106 +++++++++++++++++++++
.../table/source/snapshot/ScannerTestBase.java | 6 +-
4 files changed, 115 insertions(+), 6 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableBatchScan.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableBatchScan.java
index 710f736985..6928ad9f9b 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableBatchScan.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableBatchScan.java
@@ -177,10 +177,7 @@ public class DataTableBatchScan extends
AbstractDataTableScan {
}
private Optional<StartingScanner.Result> applyPushDownTopN() {
- if (topN == null
- || pushDownLimit != null
- || !schema.primaryKeys().isEmpty()
- || options().deletionVectorsEnabled()) {
+ if (topN == null || pushDownLimit != null ||
!schema.primaryKeys().isEmpty()) {
return Optional.empty();
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/TopNDataSplitEvaluator.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/TopNDataSplitEvaluator.java
index b33f6e9212..d6bf3d3215 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/TopNDataSplitEvaluator.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/TopNDataSplitEvaluator.java
@@ -37,6 +37,7 @@ import java.util.stream.Collectors;
import static org.apache.paimon.predicate.SortValue.NullOrdering.NULLS_FIRST;
import static org.apache.paimon.predicate.SortValue.SortDirection.ASCENDING;
import static org.apache.paimon.table.source.PushDownUtils.minmaxAvailable;
+import static
org.apache.paimon.table.source.PushDownUtils.tightBoundsAvailable;
/** Evaluate DataSplit TopN result. */
public class TopNDataSplitEvaluator {
@@ -68,7 +69,8 @@ public class TopNDataSplitEvaluator {
List<Split> results = new ArrayList<>();
List<RichSplit> richSplits = new ArrayList<>();
for (Split split : splits) {
- if (!minmaxAvailable(split, Collections.singleton(field.name()))) {
+ if (!minmaxAvailable(split, Collections.singleton(field.name()))
+ || !tightBoundsAvailable(split)) {
// unknown split, read it
results.add(split);
continue;
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java
index 90cfb7be1a..6c25574b02 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java
@@ -18,10 +18,17 @@
package org.apache.paimon.table.source;
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.BinaryRowWriter;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.manifest.FileSource;
+import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.FieldRef;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.predicate.TopN;
+import org.apache.paimon.stats.SimpleStats;
import org.apache.paimon.stats.SimpleStatsEvolutions;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.sink.StreamTableCommit;
@@ -38,6 +45,7 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
+import static org.apache.paimon.data.BinaryArray.fromLongArray;
import static org.apache.paimon.predicate.SortValue.NullOrdering.NULLS_FIRST;
import static org.apache.paimon.predicate.SortValue.NullOrdering.NULLS_LAST;
import static org.apache.paimon.predicate.SortValue.SortDirection.ASCENDING;
@@ -530,6 +538,58 @@ public class TableScanTest extends ScannerTestBase {
commit.close();
}
+ @Test
+ public void testPushDownTopNWithDeletionVectorsEnabled() throws Exception {
+ Options options = new Options();
+ options.set(CoreOptions.DELETION_VECTORS_ENABLED, true);
+ createAppendOnlyTable(options);
+
+ StreamTableWrite write = table.newWrite(commitUser);
+ StreamTableCommit commit = table.newCommit(commitUser);
+
+ for (int i = 1; i <= 5; i++) {
+ write.write(rowData(i, i * 10, i * 100L));
+ commit.commit(i, write.prepareCommit(true, i));
+ }
+ write.close();
+ commit.close();
+
+ assertThat(table.newScan().plan().splits()).hasSize(5);
+
+ DataField field = table.schema().fields().get(1);
+ FieldRef ref = new FieldRef(1, field.name(), field.type());
+ SimpleStatsEvolutions evolutions =
+ new SimpleStatsEvolutions(
+ (id) -> table.schemaManager().schema(id).fields(),
table.schema().id());
+
+ TableScan.Plan plan =
+ table.newScan().withTopN(new TopN(ref, ASCENDING, NULLS_LAST,
1)).plan();
+ assertThat(plan.splits()).hasSize(1);
+ assertThat(((DataSplit) plan.splits().get(0)).minValue(1, field,
evolutions)).isEqualTo(10);
+ }
+
+ @Test
+ public void testPushDownTopNKeepsWideDeletionVectorSplits() throws
Exception {
+ createAppendOnlyTable();
+
+ DataField field = table.schema().fields().get(1);
+ FieldRef ref = new FieldRef(1, field.name(), field.type());
+ TopN topN = new TopN(ref, ASCENDING, NULLS_LAST, 1);
+
+ DataSplit tightLowSplit = newTestSplit("tight-low", 10, 19, null);
+ DataSplit tightHighSplit = newTestSplit("tight-high", 100, 109, null);
+ DataSplit wideSplit = newTestSplit("wide", 1000, 1009, new
DeletionFile("dv", 0, 0, 1L));
+
+ List<Split> result =
+ new TopNDataSplitEvaluator(table.schema(),
table.schemaManager())
+ .evaluate(
+ topN.orders().get(0),
+ topN.limit(),
+ Arrays.asList(tightLowSplit, tightHighSplit,
wideSplit));
+
+ assertThat(result).containsExactly(wideSplit, tightLowSplit);
+ }
+
@Test
public void testPushDownTopNSchemaEvolution() throws Exception {
createAppendOnlyTable();
@@ -619,4 +679,50 @@ public class TableScanTest extends ScannerTestBase {
assertThat(((DataSplit) plan2.splits().get(0)).maxValue(field.id(),
field, evolutions))
.isNull();
}
+
+ private DataSplit newTestSplit(
+ String name, int minValue, int maxValue, DeletionFile
deletionFile) {
+ DataFileMeta file =
+ DataFileMeta.forAppend(
+ name,
+ 0,
+ maxValue - minValue + 1L,
+ new SimpleStats(
+ newStatsRow(0, minValue, 0L),
+ newStatsRow(0, maxValue, 0L),
+ fromLongArray(new Long[] {0L, 0L, 0L})),
+ 0,
+ 0,
+ table.schema().id(),
+ Collections.emptyList(),
+ null,
+ FileSource.APPEND,
+ null,
+ null,
+ null,
+ null);
+
+ DataSplit.Builder builder =
+ DataSplit.builder()
+ .withSnapshot(1)
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("dummy")
+ .rawConvertible(true)
+ .withDataFiles(Collections.singletonList(file));
+ if (deletionFile != null) {
+
builder.withDataDeletionFiles(Collections.singletonList(deletionFile));
+ }
+ return builder.build();
+ }
+
+ private BinaryRow newStatsRow(int pt, int a, long b) {
+ BinaryRow row = new BinaryRow(3);
+ BinaryRowWriter writer = new BinaryRowWriter(row);
+ writer.writeInt(0, pt);
+ writer.writeInt(1, a);
+ writer.writeLong(2, b);
+ writer.complete();
+ return row;
+ }
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ScannerTestBase.java
b/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ScannerTestBase.java
index 612b4ed49b..2f92f90565 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ScannerTestBase.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ScannerTestBase.java
@@ -87,10 +87,14 @@ public abstract class ScannerTestBase {
}
protected void createAppendOnlyTable() throws Exception {
+ createAppendOnlyTable(new Options());
+ }
+
+ protected void createAppendOnlyTable(Options conf) throws Exception {
tempDir = Files.createTempDirectory("junit");
tablePath = new Path(TraceableFileIO.SCHEME + "://" +
tempDir.toString());
fileIO = FileIOFinder.find(tablePath);
- table = createFileStoreTable(false);
+ table = createFileStoreTable(false, conf, tablePath);
snapshotReader = table.newSnapshotReader();
}