From 1905b57952d8b8d77cee4710cf5096bd5e969a40 Mon Sep 17 00:00:00 2001 From: jianguotian <18464293+jianguotian@users.noreply.github.com> Date: Sun, 13 Sep 2026 22:27:55 +0800 Subject: [PATCH] [core] Support bucket pruning for manifest files --- .../apache/paimon/manifest/BucketFilter.java | 25 +++++++++ .../paimon/manifest/ManifestAvroWriter.java | 33 +++++++++++- .../paimon/manifest/ManifestBucketFilter.java | 29 ++++++++++ .../paimon/manifest/ManifestFileMeta.java | 54 ++++++++++++++++--- .../manifest/ManifestFileMetaSerializer.java | 6 ++- .../manifest/ProjectedManifestEntry.java | 1 + .../operation/AbstractFileStoreScan.java | 2 + .../paimon/operation/BucketSelector.java | 32 +++++++++-- .../operation/ManifestEntryRunMerge.java | 19 +++++++ .../operation/ManifestEntryRunMergePlan.java | 2 + .../operation/ManifestFileBlockMerger.java | 21 ++++++++ .../paimon/operation/ManifestsReader.java | 40 +++++++++----- .../ManifestFileMetaSerializerTest.java | 3 +- .../paimon/manifest/ManifestFileTest.java | 34 ++++++++++++ .../paimon/manifest/ManifestListTest.java | 17 ++++-- .../manifest/ManifestTestDataGenerator.java | 18 ++++++- .../paimon/operation/BucketSelectorTest.java | 51 ++++++++++++++++++ 17 files changed, 358 insertions(+), 29 deletions(-) create mode 100644 paimon-core/src/main/java/org/apache/paimon/manifest/ManifestBucketFilter.java diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/BucketFilter.java b/paimon-core/src/main/java/org/apache/paimon/manifest/BucketFilter.java index 4662d7dab517..536be456eef0 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/BucketFilter.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/BucketFilter.java @@ -77,4 +77,29 @@ public boolean test(BinaryRow partition, int bucket, int totalBucket) { return totalAwareBucketFilter == null || totalAwareBucketFilter.test(partition, bucket, totalBucket); } + + /** Conservatively tests whether a manifest's bucket metadata can contain a matching entry. */ + public boolean mayContain(ManifestFileMeta manifest) { + Integer minBucket = manifest.minBucket(); + Integer maxBucket = manifest.maxBucket(); + if (minBucket == null || maxBucket == null) { + return true; + } + if (onlyReadRealBuckets && maxBucket < 0) { + return false; + } + if (specifiedBucket != null + && (specifiedBucket < minBucket || specifiedBucket > maxBucket)) { + return false; + } + if (totalAwareBucketFilter instanceof ManifestBucketFilter) { + Integer totalBuckets = manifest.totalBuckets(); + if (minBucket < 0 || totalBuckets == null || totalBuckets <= 0) { + return true; + } + return ((ManifestBucketFilter) totalAwareBucketFilter) + .mayContain(minBucket, maxBucket, totalBuckets); + } + return true; + } } diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroWriter.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroWriter.java index e78f273e29f1..c0eb377b7971 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroWriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroWriter.java @@ -296,6 +296,7 @@ public static final class EncodedEntry { private byte kind; private BinaryRow partition; private int bucket; + private int totalBuckets; private int level; private long schemaId; private boolean hasRowId; @@ -306,6 +307,7 @@ public EncodedEntry replace( byte kind, BinaryRow partition, int bucket, + int totalBuckets, int level, long schemaId, long firstRowId, @@ -313,6 +315,7 @@ public EncodedEntry replace( this.kind = kind; this.partition = partition; this.bucket = bucket; + this.totalBuckets = totalBuckets; this.level = level; this.schemaId = schemaId; this.hasRowId = true; @@ -325,12 +328,14 @@ public EncodedEntry replace( byte kind, BinaryRow partition, int bucket, + int totalBuckets, int level, long schemaId, long rowCount) { this.kind = kind; this.partition = partition; this.bucket = bucket; + this.totalBuckets = totalBuckets; this.level = level; this.schemaId = schemaId; this.hasRowId = false; @@ -348,6 +353,7 @@ public static final class EncodedBlockMeta { private final long schemaId; private final int minBucket; private final int maxBucket; + private final @Nullable Integer totalBuckets; private final int minLevel; private final int maxLevel; private final long minRowId; @@ -360,6 +366,7 @@ public EncodedBlockMeta( long schemaId, int minBucket, int maxBucket, + @Nullable Integer totalBuckets, int minLevel, int maxLevel, long minRowId, @@ -370,6 +377,7 @@ public EncodedBlockMeta( this.schemaId = schemaId; this.minBucket = minBucket; this.maxBucket = maxBucket; + this.totalBuckets = totalBuckets; this.minLevel = minLevel; this.maxLevel = maxLevel; this.minRowId = minRowId; @@ -396,9 +404,11 @@ private final class FileWriter { private long schemaId = Long.MIN_VALUE; private int minBucket = Integer.MAX_VALUE; private int maxBucket = Integer.MIN_VALUE; + private @Nullable Integer totalBuckets; private int minLevel = Integer.MAX_VALUE; private int maxLevel = Integer.MIN_VALUE; private boolean bucketStatsKnown = true; + private boolean totalBucketsKnown = true; private boolean levelStatsKnown = true; private @Nullable RowIdStats rowIdStats = new RowIdStats(); private boolean closed; @@ -473,6 +483,7 @@ private void collectStats(ManifestEntry entry) { schemaId = Math.max(schemaId, entry.file().schemaId()); minBucket = Math.min(minBucket, entry.bucket()); maxBucket = Math.max(maxBucket, entry.bucket()); + collectTotalBuckets(entry.totalBuckets()); minLevel = Math.min(minLevel, entry.level()); maxLevel = Math.max(maxLevel, entry.level()); if (rowIdStats != null) { @@ -500,6 +511,7 @@ private void collectStats(EncodedEntry entry) { schemaId = Math.max(schemaId, entry.schemaId); minBucket = Math.min(minBucket, entry.bucket); maxBucket = Math.max(maxBucket, entry.bucket); + collectTotalBuckets(entry.totalBuckets); minLevel = Math.min(minLevel, entry.level); maxLevel = Math.max(maxLevel, entry.level); if (rowIdStats != null) { @@ -517,6 +529,7 @@ private void collectStats(EncodedBlockMeta metadata) { schemaId = Math.max(schemaId, metadata.schemaId); minBucket = Math.min(minBucket, metadata.minBucket); maxBucket = Math.max(maxBucket, metadata.maxBucket); + collectTotalBuckets(metadata.totalBuckets); minLevel = Math.min(minLevel, metadata.minLevel); maxLevel = Math.max(maxLevel, metadata.maxLevel); if (rowIdStats != null) { @@ -538,6 +551,7 @@ private void collectStats(ManifestFileMeta manifest) { minBucket = Math.min(minBucket, manifest.minBucket()); maxBucket = Math.max(maxBucket, manifest.maxBucket()); } + collectTotalBuckets(manifest.totalBuckets()); if (manifest.minLevel() == null || manifest.maxLevel() == null) { levelStatsKnown = false; } else { @@ -555,6 +569,21 @@ private void collectStats(ManifestFileMeta manifest) { collectCopiedPartitionStats(manifest.partitionStats()); } + private void collectTotalBuckets(@Nullable Integer candidate) { + if (!totalBucketsKnown) { + return; + } + if (candidate == null || candidate <= 0) { + totalBucketsKnown = false; + totalBuckets = null; + } else if (totalBuckets == null) { + totalBuckets = candidate; + } else if (!totalBuckets.equals(candidate)) { + totalBucketsKnown = false; + totalBuckets = null; + } + } + private void collectCopiedPartitionStats(SimpleStats partitionStats) { collectCopiedPartitionRepresentative(partitionStats.minValues()); if (!partitionStats.maxValues().equals(partitionStats.minValues())) { @@ -709,7 +738,9 @@ private ManifestFileMeta result() { levelStatsKnown ? minLevel : null, levelStatsKnown ? maxLevel : null, rowIdStats == null ? null : rowIdStats.minRowId, - rowIdStats == null ? null : rowIdStats.maxRowId); + rowIdStats == null ? null : rowIdStats.maxRowId, + null, + totalBucketsKnown ? totalBuckets : null); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestBucketFilter.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestBucketFilter.java new file mode 100644 index 000000000000..f7f5417509fa --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestBucketFilter.java @@ -0,0 +1,29 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.manifest; + +import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.utils.TriFilter; + +/** A total-aware bucket filter which can conservatively prune manifest files. */ +public interface ManifestBucketFilter extends TriFilter { + + /** Returns whether a manifest bucket range may contain a matching bucket. */ + boolean mayContain(int minBucket, int maxBucket, int totalBuckets); +} diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMeta.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMeta.java index 2123a419a5d9..b0c7f7b13f7f 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMeta.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMeta.java @@ -62,8 +62,8 @@ public class ManifestFileMeta { new DataField( 12, "_EXTRA_FILES", - new ArrayType( - true, new VarCharType(false, Integer.MAX_VALUE))))); + new ArrayType(true, new VarCharType(false, Integer.MAX_VALUE))), + new DataField(13, "_TOTAL_BUCKETS", new IntType(true)))); private final String fileName; private final long fileSize; @@ -78,6 +78,7 @@ public class ManifestFileMeta { private final @Nullable Long minRowId; private final @Nullable Long maxRowId; private final @Nullable List extraFiles; + private final @Nullable Integer totalBuckets; public ManifestFileMeta( String fileName, @@ -105,6 +106,7 @@ public ManifestFileMeta( maxLevel, minRowId, maxRowId, + null, null); } @@ -122,6 +124,38 @@ public ManifestFileMeta( @Nullable Long minRowId, @Nullable Long maxRowId, @Nullable List extraFiles) { + this( + fileName, + fileSize, + numAddedFiles, + numDeletedFiles, + partitionStats, + schemaId, + minBucket, + maxBucket, + minLevel, + maxLevel, + minRowId, + maxRowId, + extraFiles, + null); + } + + public ManifestFileMeta( + String fileName, + long fileSize, + long numAddedFiles, + long numDeletedFiles, + SimpleStats partitionStats, + long schemaId, + @Nullable Integer minBucket, + @Nullable Integer maxBucket, + @Nullable Integer minLevel, + @Nullable Integer maxLevel, + @Nullable Long minRowId, + @Nullable Long maxRowId, + @Nullable List extraFiles, + @Nullable Integer totalBuckets) { this.fileName = fileName; this.fileSize = fileSize; this.numAddedFiles = numAddedFiles; @@ -135,6 +169,7 @@ public ManifestFileMeta( this.minRowId = minRowId; this.maxRowId = maxRowId; this.extraFiles = extraFiles; + this.totalBuckets = totalBuckets; } public String fileName() { @@ -189,6 +224,10 @@ public long schemaId() { return extraFiles; } + public @Nullable Integer totalBuckets() { + return totalBuckets; + } + @Override public boolean equals(Object o) { if (!(o instanceof ManifestFileMeta)) { @@ -207,7 +246,8 @@ public boolean equals(Object o) { && Objects.equals(maxLevel, that.maxLevel) && Objects.equals(minRowId, that.minRowId) && Objects.equals(maxRowId, that.maxRowId) - && Objects.equals(extraFiles, that.extraFiles); + && Objects.equals(extraFiles, that.extraFiles) + && Objects.equals(totalBuckets, that.totalBuckets); } @Override @@ -225,13 +265,14 @@ public int hashCode() { maxLevel, minRowId, maxRowId, - extraFiles); + extraFiles, + totalBuckets); } @Override public String toString() { return String.format( - "{%s, %d, %d, %d, %s, %d, %s, %s, %s, %s, %s, %s, %s}", + "{%s, %d, %d, %d, %s, %d, %s, %s, %s, %s, %s, %s, %s, %s}", fileName, fileSize, numAddedFiles, @@ -244,7 +285,8 @@ public String toString() { maxLevel, minRowId, maxRowId, - extraFiles); + extraFiles, + totalBuckets); } // ----------------------- Serialization ----------------------------- diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMetaSerializer.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMetaSerializer.java index 4c2ab90c96ff..7e44d42f995c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMetaSerializer.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFileMetaSerializer.java @@ -60,7 +60,8 @@ public InternalRow toRow(ManifestFileMeta meta) { meta.maxLevel(), meta.minRowId(), meta.maxRowId(), - toStringArrayData(meta.extraFiles())); + toStringArrayData(meta.extraFiles()), + meta.totalBuckets()); } @Override @@ -95,6 +96,7 @@ private ManifestFileMeta fromDataRow(InternalRow row) { row.isNullAt(9) ? null : row.getInt(9), row.isNullAt(10) ? null : row.getLong(10), row.isNullAt(11) ? null : row.getLong(11), - row.isNullAt(12) ? null : fromStringArrayData(row.getArray(12))); + row.isNullAt(12) ? null : fromStringArrayData(row.getArray(12)), + row.getFieldCount() <= 13 || row.isNullAt(13) ? null : row.getInt(13)); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java index 1d99c8bd918d..77486d4cdd53 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java @@ -142,6 +142,7 @@ private static Projection createEntryLayoutProjection() { manifestType.getField(ManifestEntry.KIND), manifestType.getField(ManifestEntry.PARTITION), manifestType.getField(ManifestEntry.BUCKET), + manifestType.getField(ManifestEntry.TOTAL_BUCKETS), manifestType .getField(ManifestEntry.FILE) .newType( diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java index a9ef5902ec9e..82923ef9d827 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java @@ -157,6 +157,7 @@ public FileStoreScan withBucket(int bucket) { @Override public FileStoreScan withBucketFilter(Filter bucketFilter) { + manifestsReader.withBucketFilter(bucketFilter); this.bucketFilter = bucketFilter; return this; } @@ -164,6 +165,7 @@ public FileStoreScan withBucketFilter(Filter bucketFilter) { @Override public FileStoreScan withTotalAwareBucketFilter( TriFilter totalAwareBucketFilter) { + manifestsReader.withTotalAwareBucketFilter(totalAwareBucketFilter); this.totalAwareBucketFilter = totalAwareBucketFilter; return this; } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/BucketSelector.java b/paimon-core/src/main/java/org/apache/paimon/operation/BucketSelector.java index 8c28c3c0c641..ee4c1c1ec23a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/BucketSelector.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/BucketSelector.java @@ -23,6 +23,7 @@ import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.GenericRow; import org.apache.paimon.data.serializer.InternalRowSerializer; +import org.apache.paimon.manifest.ManifestBucketFilter; import org.apache.paimon.predicate.Equal; import org.apache.paimon.predicate.FieldRef; import org.apache.paimon.predicate.In; @@ -31,7 +32,6 @@ import org.apache.paimon.predicate.Predicate; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.BiFilter; -import org.apache.paimon.utils.TriFilter; import org.apache.paimon.shade.guava30.com.google.common.collect.ImmutableSet; @@ -55,7 +55,7 @@ /** Selector to select bucket from {@link Predicate}. */ @ThreadSafe -public class BucketSelector implements TriFilter { +public class BucketSelector implements ManifestBucketFilter { public static final int MAX_VALUES = 1000; @@ -65,6 +65,7 @@ public class BucketSelector implements TriFilter { private final RowType bucketKeyType; private final Predicate predicate; private final Map> partitionSelectors; + private final Optional manifestSelector; public BucketSelector( Predicate predicate, @@ -78,6 +79,7 @@ public BucketSelector( this.partitionType = partitionType; this.bucketKeyType = bucketKeyType; this.partitionSelectors = new ConcurrentHashMap<>(); + this.manifestSelector = createPartitionSelectorFromPredicate(predicate); } @Override @@ -88,6 +90,16 @@ public boolean test(BinaryRow partition, Integer bucket, Integer numBucket) { .orElse(true); } + @Override + public boolean mayContain(int minBucket, int maxBucket, int totalBuckets) { + if (minBucket < 0 || maxBucket < minBucket || totalBuckets <= 0) { + return true; + } + return manifestSelector + .map(selector -> selector.mayContain(minBucket, maxBucket, totalBuckets)) + .orElse(true); + } + private Optional createPartitionSelector(BinaryRow partition) { Optional partRemoved = predicate.visit( @@ -96,9 +108,14 @@ private Optional createPartitionSelector(BinaryRow partition) return Optional.empty(); } + return createPartitionSelectorFromPredicate(partRemoved.get()); + } + + private Optional createPartitionSelectorFromPredicate( + Predicate sourcePredicate) { List bucketFilters = pickTransformFieldMapping( - splitAnd(partRemoved.get()), + splitAnd(sourcePredicate), rowType.getFieldNames(), bucketKeyType.getFieldNames()); if (bucketFilters.isEmpty()) { @@ -226,5 +243,14 @@ private Set createBucketSet(int numBucket) { } return builder.build(); } + + private boolean mayContain(int minBucket, int maxBucket, int totalBuckets) { + for (Integer bucket : buckets.computeIfAbsent(totalBuckets, this::createBucketSet)) { + if (bucket >= minBucket && bucket <= maxBucket) { + return true; + } + } + return false; + } } } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java index 7eb4f66e32c8..740b23f58d6d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java @@ -630,6 +630,8 @@ private static final class BlockMetadataAccumulator { private long schemaId = Long.MIN_VALUE; private int minBucket = Integer.MAX_VALUE; private int maxBucket = Integer.MIN_VALUE; + private @Nullable Integer totalBuckets; + private boolean totalBucketsKnown = true; private int minLevel = Integer.MAX_VALUE; private int maxLevel = Integer.MIN_VALUE; private long minRowId = Long.MAX_VALUE; @@ -650,6 +652,7 @@ private void collect(ProjectedManifestEntry entry, SortKey key, BinaryRow partit int bucket = entry.bucket(); minBucket = Math.min(minBucket, bucket); maxBucket = Math.max(maxBucket, bucket); + collectTotalBuckets(entry.totalBuckets()); int level = entry.file().level(); minLevel = Math.min(minLevel, level); maxLevel = Math.max(maxLevel, level); @@ -665,12 +668,28 @@ private EncodedBlockMeta finish(SimpleStatsConverter partitionStatsConverter) { schemaId, minBucket, maxBucket, + totalBucketsKnown ? totalBuckets : null, minLevel, maxLevel, minRowId, maxRowId, partitionStats.finish(partitionStatsConverter)); } + + private void collectTotalBuckets(int candidate) { + if (!totalBucketsKnown) { + return; + } + if (candidate <= 0) { + totalBucketsKnown = false; + totalBuckets = null; + } else if (totalBuckets == null) { + totalBuckets = candidate; + } else if (totalBuckets != candidate) { + totalBucketsKnown = false; + totalBuckets = null; + } + } } /** Partition statistics for one sorted Avro block. */ diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java index 618943df180d..1645123f9128 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java @@ -483,6 +483,7 @@ public boolean advance() throws Exception { key.kind, partitions.partition(key.partitionId), currentEntry.bucket(), + currentEntry.totalBuckets(), currentEntry.file().level(), currentEntry.file().schemaId(), key.firstRowId, @@ -636,6 +637,7 @@ public void materializeCurrent() throws Exception { key.kind, partitions.partition(key.partitionId), currentEntry.bucket(), + currentEntry.totalBuckets(), currentEntry.file().level(), currentEntry.file().schemaId(), key.firstRowId, diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileBlockMerger.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileBlockMerger.java index a3631122156a..04a4a7d516f9 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileBlockMerger.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileBlockMerger.java @@ -742,6 +742,7 @@ private static void writeCompactedEntry( entry.kind().toByteValue(), partition, entry.bucket(), + entry.totalBuckets(), file.level(), file.schemaId(), file.nonNullFirstRowId(), @@ -751,6 +752,7 @@ private static void writeCompactedEntry( entry.kind().toByteValue(), partition, entry.bucket(), + entry.totalBuckets(), file.level(), file.schemaId(), file.rowCount()); @@ -771,6 +773,8 @@ private static final class CompactionBlock { private long schemaId = Long.MIN_VALUE; private int minBucket = Integer.MAX_VALUE; private int maxBucket = Integer.MIN_VALUE; + private @Nullable Integer totalBuckets; + private boolean totalBucketsKnown = true; private int minLevel = Integer.MAX_VALUE; private int maxLevel = Integer.MIN_VALUE; private long minRowId = Long.MAX_VALUE; @@ -817,6 +821,7 @@ private boolean collect( int bucket = entry.bucket(); minBucket = Math.min(minBucket, bucket); maxBucket = Math.max(maxBucket, bucket); + collectTotalBuckets(entry.totalBuckets()); int level = file.level(); minLevel = Math.min(minLevel, level); maxLevel = Math.max(maxLevel, level); @@ -864,6 +869,7 @@ private void finish( schemaId, minBucket, maxBucket, + totalBucketsKnown ? totalBuckets : null, minLevel, maxLevel, hasRowIds ? minRowId : -1, @@ -872,6 +878,21 @@ private void finish( partitionCounts = null; } + private void collectTotalBuckets(int candidate) { + if (!totalBucketsKnown) { + return; + } + if (candidate <= 0) { + totalBucketsKnown = false; + totalBuckets = null; + } else if (totalBuckets == null) { + totalBuckets = candidate; + } else if (totalBuckets != candidate) { + totalBucketsKnown = false; + totalBuckets = null; + } + } + private void finishFiltering(CollectedDeletes deletes, boolean deferDeletedAddCheck) { if (deferDeletedAddCheck && metadata != null) { checkState(hasRowIds, "RowID filtering requires block RowID statistics."); diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestsReader.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestsReader.java index 46bbbd17f6e3..618eae930d40 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestsReader.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestsReader.java @@ -20,6 +20,7 @@ import org.apache.paimon.Snapshot; import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.manifest.BucketFilter; import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.manifest.ManifestList; import org.apache.paimon.partition.PartitionPredicate; @@ -28,8 +29,10 @@ import org.apache.paimon.table.source.ScanMode; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.BiFilter; +import org.apache.paimon.utils.Filter; import org.apache.paimon.utils.RowRangeIndex; import org.apache.paimon.utils.SnapshotManager; +import org.apache.paimon.utils.TriFilter; import javax.annotation.Nullable; import javax.annotation.concurrent.ThreadSafe; @@ -53,6 +56,8 @@ public class ManifestsReader { private boolean onlyReadRealBuckets = false; @Nullable private Integer specifiedBucket = null; + @Nullable private Filter bucketFilter = null; + @Nullable private TriFilter totalAwareBucketFilter = null; @Nullable private Integer specifiedLevel = null; @Nullable private PartitionPredicate partitionFilter = null; // Auth partition filter (ANDed with partitionFilter); kept separate so it can be reset each @@ -82,6 +87,17 @@ public ManifestsReader withBucket(int bucket) { return this; } + public ManifestsReader withBucketFilter(Filter bucketFilter) { + this.bucketFilter = bucketFilter; + return this; + } + + public ManifestsReader withTotalAwareBucketFilter( + TriFilter totalAwareBucketFilter) { + this.totalAwareBucketFilter = totalAwareBucketFilter; + return this; + } + public ManifestsReader withLevel(int level) { this.specifiedLevel = level; return this; @@ -147,9 +163,15 @@ public Result read(@Nullable Snapshot specifiedSnapshot, ScanMode scanMode) { // Compute the effective partition filter once (it ANDs the base and auth slots) instead of // rebuilding it per manifest. PartitionPredicate effectivePartitionFilter = partitionFilter(); + BucketFilter effectiveBucketFilter = + BucketFilter.create( + onlyReadRealBuckets, specifiedBucket, bucketFilter, totalAwareBucketFilter); List filtered = manifests.stream() - .filter(m -> filterManifestFileMeta(m, effectivePartitionFilter)) + .filter( + m -> + filterManifestFileMeta( + m, effectivePartitionFilter, effectiveBucketFilter)) .collect(Collectors.toList()); return new Result(snapshot, manifests, filtered); } @@ -183,17 +205,11 @@ private boolean filterManifestByRowRanges(ManifestFileMeta manifest) { /** Note: Keep this thread-safe. */ private boolean filterManifestFileMeta( - ManifestFileMeta manifest, @Nullable PartitionPredicate effectivePartitionFilter) { - Integer minBucket = manifest.minBucket(); - Integer maxBucket = manifest.maxBucket(); - if (minBucket != null && maxBucket != null) { - if (onlyReadRealBuckets && maxBucket < 0) { - return false; - } - if (specifiedBucket != null - && (specifiedBucket < minBucket || specifiedBucket > maxBucket)) { - return false; - } + ManifestFileMeta manifest, + @Nullable PartitionPredicate effectivePartitionFilter, + @Nullable BucketFilter effectiveBucketFilter) { + if (effectiveBucketFilter != null && !effectiveBucketFilter.mayContain(manifest)) { + return false; } Integer minLevel = manifest.minLevel(); diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaSerializerTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaSerializerTest.java index 2f4e32cd2759..7600c7d01fb7 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaSerializerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileMetaSerializerTest.java @@ -68,7 +68,8 @@ void testExtraFiles() throws IOException { original.maxLevel(), original.minRowId(), original.maxRowId(), - extraFiles); + extraFiles, + original.totalBuckets()); ManifestFileMeta fromRow = serializer.fromRow(serializer.toRow(meta)); ManifestFileMeta fromBytes = serializer.deserializeFromBytes(meta.toBytes()); diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java index c3a50f4ef1de..c472eb1d3411 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java @@ -193,6 +193,7 @@ void testWriteEncodedRecords() throws Exception { entry.kind().toByteValue(), entry.partition(), entry.bucket(), + entry.totalBuckets(), entry.level(), entry.file().schemaId(), entry.file().firstRowId(), @@ -210,6 +211,7 @@ void testWriteEncodedRecords() throws Exception { assertThat(result.schemaId()).isEqualTo(sourceMeta.schemaId()); assertThat(result.minBucket()).isEqualTo(sourceMeta.minBucket()); assertThat(result.maxBucket()).isEqualTo(sourceMeta.maxBucket()); + assertThat(result.totalBuckets()).isEqualTo(sourceMeta.totalBuckets()); assertThat(result.minLevel()).isEqualTo(sourceMeta.minLevel()); assertThat(result.maxLevel()).isEqualTo(sourceMeta.maxLevel()); assertThat(result.minRowId()).isEqualTo(sourceMeta.minRowId()); @@ -246,6 +248,7 @@ void testWriteEncodedRecordsFlushesPartitionStatsBuffer() throws Exception { source.kind().toByteValue(), source.partition().copy(), source.bucket(), + source.totalBuckets(), source.level(), source.file().schemaId(), source.file().firstRowId(), @@ -289,12 +292,42 @@ void testWriteEncodedManifestPreservesUnknownAggregateStats() throws Exception { ManifestFileMeta result = writer.result().get(0); assertThat(result.minBucket()).isNull(); assertThat(result.maxBucket()).isNull(); + assertThat(result.totalBuckets()).isNull(); assertThat(result.minLevel()).isNull(); assertThat(result.maxLevel()).isNull(); assertThat(result.partitionStats()).isEqualTo(source.partitionStats()); assertThat(manifestFile.read(result.fileName())).containsExactlyElementsOf(entries); } + @Test + void testTotalBucketsAggregateStats() throws Exception { + ManifestEntry source = gen.next(); + ManifestFile manifestFile = createManifestFile(tempDir.toString(), Long.MAX_VALUE); + + ManifestEntry add = + ManifestEntry.create( + FileKind.ADD, source.partition(), source.bucket(), 8, source.file()); + ManifestEntry delete = + ManifestEntry.create( + FileKind.DELETE, source.partition(), source.bucket(), 8, source.file()); + assertThat(writeSingleManifest(manifestFile, Arrays.asList(add, delete)).totalBuckets()) + .isEqualTo(8); + + ManifestEntry different = + ManifestEntry.create( + FileKind.ADD, source.partition(), source.bucket(), 16, source.file()); + assertThat(writeSingleManifest(manifestFile, Arrays.asList(add, different)).totalBuckets()) + .isNull(); + + ManifestEntry nonPositive = + ManifestEntry.create( + FileKind.DELETE, source.partition(), source.bucket(), 0, source.file()); + assertThat( + writeSingleManifest(manifestFile, Arrays.asList(add, nonPositive)) + .totalBuckets()) + .isNull(); + } + @Test void testReadMissingManifestFile() { ManifestFile manifestFile = createManifestFile(tempDir.toString()); @@ -1185,6 +1218,7 @@ private ManifestAvroWriter.EncodedBlockMeta encodedBlockMeta(ManifestFileMeta me meta.schemaId(), meta.minBucket(), meta.maxBucket(), + meta.totalBuckets(), meta.minLevel(), meta.maxLevel(), meta.minRowId() == null ? -1 : meta.minRowId(), diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestListTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestListTest.java index 4891a760b7ff..37b840b9fe4a 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestListTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestListTest.java @@ -101,7 +101,12 @@ public void testCanReadOldMetaPaimon10() throws Exception { ManifestList manifestList = createManifestList(tempDir.toString()); List actualMetas = manifestList.read(manifestListName); assertThat(actualMetas).isEqualTo(getLegacyMetaPaimon10(metas)); - assertThat(actualMetas).allSatisfy(meta -> assertThat(meta.extraFiles()).isNull()); + assertThat(actualMetas) + .allSatisfy( + meta -> { + assertThat(meta.extraFiles()).isNull(); + assertThat(meta.totalBuckets()).isNull(); + }); } @Test @@ -113,7 +118,12 @@ public void testOldReaderCanReadNewMetaPaimon10() throws Exception { ManifestList legacyManifestList = createLegacyManifestListPaimon10(); List actualMetas = legacyManifestList.read(manifestListName); assertThat(actualMetas).isEqualTo(getLegacyMetaPaimon10(metas)); - assertThat(actualMetas).allSatisfy(meta -> assertThat(meta.extraFiles()).isNull()); + assertThat(actualMetas) + .allSatisfy( + meta -> { + assertThat(meta.extraFiles()).isNull(); + assertThat(meta.totalBuckets()).isNull(); + }); } private ManifestList createLegacyManifestListPaimon10() { @@ -184,7 +194,8 @@ private List generateData() { meta.maxLevel(), meta.minRowId(), meta.maxRowId(), - extraFiles)); + extraFiles, + meta.totalBuckets())); } return metas; } diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestTestDataGenerator.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestTestDataGenerator.java index 4b576c6bd64c..a9cbb498e9ee 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestTestDataGenerator.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestTestDataGenerator.java @@ -101,6 +101,8 @@ public ManifestFileMeta createManifestFileMeta(List entries) { long numDeletedFiles = 0; int minBucket = Integer.MAX_VALUE; int maxBucket = Integer.MIN_VALUE; + Integer totalBuckets = null; + boolean totalBucketsKnown = true; int minLevel = Integer.MAX_VALUE; int maxLevel = Integer.MIN_VALUE; for (ManifestEntry entry : entries) { @@ -112,6 +114,18 @@ public ManifestFileMeta createManifestFileMeta(List entries) { } minBucket = Math.min(minBucket, entry.bucket()); maxBucket = Math.max(maxBucket, entry.bucket()); + int candidate = entry.totalBuckets(); + if (totalBucketsKnown) { + if (candidate <= 0) { + totalBucketsKnown = false; + totalBuckets = null; + } else if (totalBuckets == null) { + totalBuckets = candidate; + } else if (totalBuckets != candidate) { + totalBucketsKnown = false; + totalBuckets = null; + } + } minLevel = Math.min(minLevel, entry.level()); maxLevel = Math.max(maxLevel, entry.level()); } @@ -128,7 +142,9 @@ public ManifestFileMeta createManifestFileMeta(List entries) { minLevel, maxLevel, null, - null); + null, + null, + totalBucketsKnown ? totalBuckets : null); } private void mergeLevelsIfNeeded(BinaryRow partition, int bucket) { diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/BucketSelectorTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/BucketSelectorTest.java index 3128ff20db22..bcee6762edc6 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/BucketSelectorTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/BucketSelectorTest.java @@ -20,6 +20,8 @@ import org.apache.paimon.CoreOptions.BucketFunctionType; import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.manifest.BucketFilter; +import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.predicate.Predicate; import org.apache.paimon.predicate.PredicateBuilder; import org.apache.paimon.types.DataTypes; @@ -31,6 +33,7 @@ import java.util.HashSet; import java.util.Set; +import static org.apache.paimon.stats.SimpleStats.EMPTY_STATS; import static org.assertj.core.api.Assertions.assertThat; /** Tests for {@link BucketSelector}. */ @@ -57,6 +60,54 @@ public void testEqualPredicate() { assertThat(selected).hasSize(1); } + @Test + public void testManifestBucketRange() { + RowType rowType = DataTypes.ROW(DataTypes.FIELD(0, "k", DataTypes.INT())); + RowType partType = RowType.of(); + RowType bucketKeyType = DataTypes.ROW(DataTypes.FIELD(0, "k", DataTypes.INT())); + PredicateBuilder pb = new PredicateBuilder(rowType); + BucketSelector selector = + new BucketSelector( + pb.equal(0, 5), + BucketFunctionType.DEFAULT, + rowType, + partType, + bucketKeyType); + + int selected = + selectedBuckets(selector, BinaryRow.EMPTY_ROW, NUM_BUCKETS).iterator().next(); + assertThat(selector.mayContain(selected, selected, NUM_BUCKETS)).isTrue(); + int different = (selected + 1) % NUM_BUCKETS; + assertThat(selector.mayContain(different, different, NUM_BUCKETS)).isFalse(); + + assertThat(selector.mayContain(-1, selected, NUM_BUCKETS)).isTrue(); + assertThat(selector.mayContain(0, 0, 0)).isTrue(); + + BucketFilter filter = new BucketFilter(false, null, null, selector); + assertThat(filter.mayContain(manifest(selected, selected, NUM_BUCKETS))).isTrue(); + assertThat(filter.mayContain(manifest(different, different, NUM_BUCKETS))).isFalse(); + assertThat(filter.mayContain(manifest(different, different, null))).isTrue(); + assertThat(filter.mayContain(manifest(-1, different, NUM_BUCKETS))).isTrue(); + } + + private static ManifestFileMeta manifest(int minBucket, int maxBucket, Integer totalBuckets) { + return new ManifestFileMeta( + "manifest", + 1, + 1, + 0, + EMPTY_STATS, + 0, + minBucket, + maxBucket, + 0, + 0, + null, + null, + null, + totalBuckets); + } + @Test public void testEqualAndRangePredicate() { // k = 5 AND k < 100 => should still select bucket for k=5