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 8335a3cefa [core] Support file index and predicate push down in 
DataEvolutionSplitRead (#8839)
8335a3cefa is described below

commit 8335a3cefa6f7e2a84835e17779dc7b271b8f86f
Author: Vova Kolmakov <[email protected]>
AuthorDate: Thu Jul 30 09:36:34 2026 +0700

    [core] Support file index and predicate push down in DataEvolutionSplitRead 
(#8839)
---
 .../paimon/operation/DataEvolutionSplitRead.java   | 332 ++++++++-
 .../paimon/table/DataEvolutionFileIndexTest.java   | 805 +++++++++++++++++++++
 2 files changed, 1107 insertions(+), 30 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
index 1941dfd85d..e5bc953109 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
@@ -26,6 +26,9 @@ import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.deletionvectors.ApplyDeletionVectorReader;
 import org.apache.paimon.deletionvectors.DeletionVector;
 import org.apache.paimon.disk.IOManager;
+import org.apache.paimon.fileindex.FileIndexResult;
+import org.apache.paimon.fileindex.bitmap.ApplyBitmapIndexRecordReader;
+import org.apache.paimon.fileindex.bitmap.BitmapIndexResult;
 import org.apache.paimon.format.FileFormatDiscover;
 import org.apache.paimon.format.FormatKey;
 import org.apache.paimon.format.FormatReaderContext;
@@ -36,13 +39,16 @@ import 
org.apache.paimon.globalindex.IndexedSplitRecordReader;
 import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.io.DataFilePathFactory;
 import org.apache.paimon.io.DataFileRecordReader;
+import org.apache.paimon.io.FileIndexEvaluator;
 import org.apache.paimon.mergetree.compact.ConcatRecordReader;
 import org.apache.paimon.partition.PartitionUtils;
 import org.apache.paimon.predicate.Predicate;
 import org.apache.paimon.reader.DataEvolutionFileReader;
+import org.apache.paimon.reader.EmptyFileRecordReader;
 import org.apache.paimon.reader.FileRecordReader;
 import org.apache.paimon.reader.ReaderSupplier;
 import org.apache.paimon.reader.RecordReader;
+import org.apache.paimon.schema.SchemaEvolutionUtil;
 import org.apache.paimon.schema.SchemaManager;
 import org.apache.paimon.schema.TableSchema;
 import org.apache.paimon.table.SpecialFields;
@@ -66,9 +72,11 @@ import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Set;
 import java.util.TreeMap;
 import java.util.function.Function;
 import java.util.function.ToLongFunction;
@@ -78,17 +86,27 @@ import static java.lang.String.format;
 import static java.util.Collections.reverseOrder;
 import static java.util.Comparator.comparingLong;
 import static org.apache.paimon.format.blob.BlobFileFormat.isBlobFile;
+import static 
org.apache.paimon.predicate.PredicateBuilder.excludePredicateWithFields;
+import static org.apache.paimon.predicate.PredicateBuilder.splitAnd;
+import static org.apache.paimon.predicate.PredicateVisitor.collectFieldNames;
 import static org.apache.paimon.table.SpecialFields.rowTypeWithRowTracking;
 import static org.apache.paimon.types.BlobType.isBlobFileField;
 import static org.apache.paimon.types.VectorType.isVectorStoreFile;
 import static org.apache.paimon.utils.DataEvolutionUtils.retrieveAnchorFile;
+import static org.apache.paimon.utils.ListUtils.isNullOrEmpty;
 import static org.apache.paimon.utils.Preconditions.checkArgument;
 import static org.apache.paimon.utils.Preconditions.checkNotNull;
 
 /**
- * A union {@link SplitRead} to read multiple inner files to merge columns, 
note that this class
- * does not support filtering push down and deletion vectors, as they can 
interfere with the process
- * of merging columns.
+ * A union {@link SplitRead} to read multiple inner files to merge columns.
+ *
+ * <p>Filters can only be pushed down where they can not interfere with the 
column merging: a file
+ * read without merging gets both the file index and the format level push 
down, while a merged
+ * group only uses the file index to skip the whole group, as dropping rows in 
one of the merged
+ * readers would break the positional alignment between them.
+ *
+ * <p>Only a filter whose every field belongs to the read type is pushed down, 
see {@link
+ * #readTypeFilters}, and only to the files that wrote those fields, see 
{@link #fileFilters}.
  *
  * <p>TODO: Optimize implementation of this class.
  */
@@ -101,10 +119,15 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
     private final FileFormatDiscover formatDiscover;
     private final FileStorePathFactory pathFactory;
     private final Map<FormatKey, FormatReaderMapping> formatReaderMappings;
+    // Kept apart from formatReaderMappings: the single file path pushes per 
file filters into the
+    // mapping, so it must not share entries with the merge path, which pushes 
none.
+    private final Map<SingleFileKey, FormatReaderMapping> 
singleFileReaderMappings;
     private final Function<Long, TableSchema> schemaFetcher;
     private final CoreOptions coreOptions;
+    private final boolean fileIndexReadEnabled;
 
     protected RowType readRowType;
+    @Nullable private List<Predicate> filters;
 
     public DataEvolutionSplitRead(
             FileIO fileIO,
@@ -122,6 +145,8 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
         this.coreOptions = coreOptions;
         this.pathFactory = pathFactory;
         this.formatReaderMappings = new HashMap<>();
+        this.singleFileReaderMappings = new HashMap<>();
+        this.fileIndexReadEnabled = coreOptions.fileIndexReadEnabled();
         this.readRowType = rowType;
     }
 
@@ -143,11 +168,29 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
 
     @Override
     public SplitRead<InternalRow> withFilter(@Nullable Predicate predicate) {
-        // TODO: Support File index push down (all conditions) and Predicate 
push down (only if no
-        // column merge)
+        if (predicate != null) {
+            this.filters = pushDownFilters(splitAnd(predicate));
+            // the single file mappings carry the filters they were built 
with, and a read can be
+            // reconfigured after it created readers, see 
AppendTableRead#innerWithFilter
+            singleFileReaderMappings.clear();
+        }
         return this;
     }
 
+    /**
+     * Row tracking fields are assigned from the manifest entry instead of 
being read from the file,
+     * and data evolution may reassign row ids, so a physical copy in the file 
can be stale. Never
+     * push them down.
+     */
+    private static List<Predicate> pushDownFilters(List<Predicate> filters) {
+        return filters.stream()
+                .filter(
+                        filter ->
+                                collectFieldNames(filter).stream()
+                                        
.noneMatch(SpecialFields::isSystemField))
+                .collect(Collectors.toList());
+    }
+
     @Override
     public RecordReader<InternalRow> createReader(Split split) throws 
IOException {
         if (split instanceof DataSplit) {
@@ -176,17 +219,12 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
                 pathFactory.createDataFilePathFactory(partition, 
dataSplit.bucket());
         List<ReaderSupplier<InternalRow>> suppliers = new ArrayList<>();
 
-        Builder formatBuilder =
-                new Builder(
-                        formatDiscover,
-                        readRowType.getFields(),
-                        // file has no row id and sequence number, they are in 
manifest entry
-                        schema ->
-                                
rowTypeWithRowTracking(schema.logicalRowType(), true, true)
-                                        .getFields(),
-                        null,
-                        null,
-                        null);
+        // the merge path builds its readers with filter push down disabled, 
so the shared builder
+        // carries no filters, the single file path creates its own with per 
file filters
+        Builder formatBuilder = formatBuilder(readRowType, null);
+        // the suppliers below run lazily, so take the filters now, the same 
way the read type is
+        // already taken by the caller
+        List<Predicate> filters = readTypeFilters(this.filters, readRowType);
 
         List<List<DataFileMeta>> splitByRowId = mergeRangesAndSort(files);
         for (List<DataFileMeta> needMergeFiles : splitByRowId) {
@@ -200,7 +238,7 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
                                     partition,
                                     dataFilePathFactory,
                                     needMergeFiles.get(0),
-                                    formatBuilder,
+                                    filters,
                                     rowRanges,
                                     readRowType,
                                     deletionVector);
@@ -209,6 +247,9 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
             } else {
                 suppliers.add(
                         () -> {
+                            if (skipByFileIndex(filters, needMergeFiles, 
dataFilePathFactory)) {
+                                return new EmptyFileRecordReader<>();
+                            }
                             DeletionVectorWithRange deletionVector =
                                     readDeletionVector(needMergeFiles, 
deletionVectorFactory);
                             return createUnionReader(
@@ -441,7 +482,8 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
                                             
DataFilePathFactory.formatIdentifier(file.fileName()),
                                             dataFilePathFactory.toPath(file),
                                             file.fileSize()),
-                                    deletionVector));
+                                    deletionVector,
+                                    null));
         }
         return ConcatRecordReader.create(readerSuppliers);
     }
@@ -459,7 +501,7 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
             BinaryRow partition,
             DataFilePathFactory dataFilePathFactory,
             DataFileMeta file,
-            Builder formatBuilder,
+            @Nullable List<Predicate> filters,
             List<Range> rowRanges,
             RowType readRowType,
             @Nullable DeletionVectorWithRange deletionVector)
@@ -467,16 +509,36 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
         FileReadTarget readTarget = readTarget(file, dataFilePathFactory, 
rowRanges);
         String formatIdentifier = readTarget.formatIdentifier;
         long schemaId = file.schemaId();
+        TableSchema dataSchema = schemaId == schema.id() ? schema : 
schemaFetcher.apply(schemaId);
+
+        // no column merge here, so the filters this file can answer reach 
both the file index and
+        // the format reader
+        List<Predicate> fileFilters = fileFilters(filters, file);
         FormatReaderMapping formatReaderMapping =
-                formatReaderMappings.computeIfAbsent(
-                        new FormatKey(file.schemaId(), formatIdentifier),
+                singleFileReaderMappings.computeIfAbsent(
+                        new SingleFileKey(
+                                schemaId, formatIdentifier, file.writeCols(), 
readRowType),
                         key ->
-                                formatBuilder.build(
-                                        formatIdentifier,
-                                        schema,
-                                        schemaId == schema.id()
-                                                ? schema
-                                                : 
schemaFetcher.apply(schemaId)));
+                                formatBuilder(readRowType, fileFilters)
+                                        .build(formatIdentifier, schema, 
dataSchema));
+
+        FileIndexResult fileIndexResult = null;
+        if (fileIndexReadEnabled) {
+            fileIndexResult =
+                    FileIndexEvaluator.evaluate(
+                            fileIO,
+                            dataSchema,
+                            devolveFilters(fileFilters, dataSchema),
+                            null,
+                            null,
+                            dataFilePathFactory,
+                            file,
+                            null);
+            if (!fileIndexResult.remain()) {
+                return new EmptyFileRecordReader<>();
+            }
+        }
+
         return createFileReader(
                 partition,
                 file,
@@ -484,7 +546,8 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
                 rowRanges,
                 readRowType,
                 readTarget,
-                deletionVector);
+                deletionVector,
+                fileIndexResult);
     }
 
     private FileRecordReader<InternalRow> createFileReader(
@@ -503,7 +566,8 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
                 rowRanges,
                 readRowType,
                 readTarget(file, dataFilePathFactory, rowRanges),
-                deletionVector);
+                deletionVector,
+                null);
     }
 
     private FileRecordReader<InternalRow> createFileReader(
@@ -513,9 +577,25 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
             List<Range> rowRanges,
             RowType readRowType,
             FileReadTarget readTarget,
-            @Nullable DeletionVectorWithRange deletionVector)
+            @Nullable DeletionVectorWithRange deletionVector,
+            @Nullable FileIndexResult fileIndexResult)
             throws IOException {
         RoaringBitmap32 selection = file.toFileSelection(rowRanges);
+        BitmapIndexResult bitmapIndexResult =
+                fileIndexResult instanceof BitmapIndexResult
+                        ? (BitmapIndexResult) fileIndexResult
+                        : null;
+        if (bitmapIndexResult != null) {
+            RoaringBitmap32 indexSelection = bitmapIndexResult.get();
+            selection =
+                    selection == null
+                            ? indexSelection.clone()
+                            : RoaringBitmap32.and(selection, indexSelection);
+            if (selection.isEmpty()) {
+                return new EmptyFileRecordReader<>();
+            }
+        }
+
         FormatReaderContext formatReaderContext =
                 new FormatReaderContext(fileIO, readTarget.path, 
readTarget.fileSize, selection);
         FileRecordReader<InternalRow> fileRecordReader =
@@ -532,6 +612,12 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
                         file.firstRowId(),
                         file.maxSequenceNumber(),
                         formatReaderMapping.getSystemFields());
+
+        if (bitmapIndexResult != null) {
+            fileRecordReader =
+                    new ApplyBitmapIndexRecordReader(fileRecordReader, 
bitmapIndexResult);
+        }
+
         return applyDeletionVector(fileRecordReader, file.nonNullRowIdRange(), 
deletionVector);
     }
 
@@ -557,6 +643,146 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
                 readerRange.from - deletionVector.range.from);
     }
 
+    /**
+     * Whether the file index proves that no row of a merged group can match 
the filters. Only plain
+     * data files are considered: {@link #mergeRangesAndSort} guarantees they 
all span the row id
+     * range of the whole group, while a blob or vector-store file only covers 
a sub range and can
+     * not prove anything for the other rows.
+     *
+     * <p>A column can be written by several files of the group; the column 
merge takes each field
+     * from the newest file that wrote it and older copies are dead. Files 
arrive newest first
+     * ({@link #mergeRangesAndSort}), so each file's index is evaluated only 
over the columns it is
+     * the newest writer of. A stale value in an older file must not veto a 
group whose winning file
+     * matches, mirroring the winner selection in {@link 
DataEvolutionFileStoreScan#evolutionStats}.
+     */
+    private boolean skipByFileIndex(
+            @Nullable List<Predicate> filters,
+            List<DataFileMeta> files,
+            DataFilePathFactory pathFactory)
+            throws IOException {
+        if (!fileIndexReadEnabled || isNullOrEmpty(filters)) {
+            return false;
+        }
+
+        Set<Integer> claimedFieldIds = new HashSet<>();
+        for (DataFileMeta file : files) {
+            if (isBlobFile(file.fileName()) || 
isVectorStoreFile(file.fileName())) {
+                continue;
+            }
+
+            TableSchema dataSchema = 
schemaFetcher.apply(file.schemaId()).project(file.writeCols());
+            // columns this file wrote but a newer file already won: their 
values here are stale
+            Set<String> overwrittenCols = new HashSet<>();
+            for (DataField field : dataSchema.fields()) {
+                if (!claimedFieldIds.add(field.id())) {
+                    overwrittenCols.add(field.name());
+                }
+            }
+
+            List<Predicate> dataFilters =
+                    excludePredicateWithFields(
+                            devolveFilters(filters, dataSchema), 
overwrittenCols);
+            if (dataFilters.isEmpty()) {
+                continue;
+            }
+
+            FileIndexResult result =
+                    FileIndexEvaluator.evaluate(
+                            fileIO, dataSchema, dataFilters, null, null, 
pathFactory, file, null);
+            if (!result.remain()) {
+                return true;
+            }
+        }
+        return false;
+    }
+
+    private Builder formatBuilder(RowType readRowType, @Nullable 
List<Predicate> filters) {
+        return new Builder(
+                formatDiscover,
+                readRowType.getFields(),
+                // file has no row id and sequence number, they are in 
manifest entry
+                schema -> rowTypeWithRowTracking(schema.logicalRowType(), 
true, true).getFields(),
+                filters,
+                null,
+                null);
+    }
+
+    /**
+     * Filters the file can answer. A column missing from a data evolution 
file is not null, its
+     * values live in another file of the same row id range, so a predicate on 
it must not be pushed
+     * down: formats such as Parquet read a column absent from the file schema 
as all null and would
+     * drop every row.
+     */
+    @Nullable
+    private List<Predicate> fileFilters(@Nullable List<Predicate> filters, 
DataFileMeta file) {
+        if (isNullOrEmpty(filters)) {
+            return null;
+        }
+
+        Set<Integer> fileFieldIds = new HashSet<>();
+        for (DataField field :
+                
schemaFetcher.apply(file.schemaId()).project(file.writeCols()).fields()) {
+            fileFieldIds.add(field.id());
+        }
+        Set<String> written = new HashSet<>();
+        for (DataField field : schema.fields()) {
+            if (fileFieldIds.contains(field.id())) {
+                written.add(field.name());
+            }
+        }
+        return filtersWithin(filters, written);
+    }
+
+    /**
+     * Filters that may be pushed down at all, those whose every field belongs 
to the read type.
+     *
+     * <p>A format reader can not evaluate the others: it only reads the 
columns it was asked for
+     * and treats the rest as absent, parquet for instance evaluates its 
column index over the
+     * requested columns only and would drop the whole row group instead of 
nothing.
+     *
+     * <p>The file index can not use them either, for a different reason. 
{@link
+     * DataEvolutionFileStoreScan#pruneByReadType} drops the files of a row id 
group that write no
+     * column of the read type, so for a column outside it the file that wins 
the column merge may
+     * not be part of the split at all. What is left is an older file whose 
copy of the column is
+     * dead, and its index would veto rows that do match. For a column of the 
read type every file
+     * writing it is kept, so the winner is always there. This holds as long 
as the read type is the
+     * one the split was planned with, or a projection of it.
+     */
+    @Nullable
+    private static List<Predicate> readTypeFilters(
+            @Nullable List<Predicate> filters, RowType readRowType) {
+        return filtersWithin(filters, new 
HashSet<>(readRowType.getFieldNames()));
+    }
+
+    /**
+     * Filters every field of which belongs to {@code fields}. A predicate is 
kept only when all of
+     * its fields qualify, a leaf can read more than one of them.
+     */
+    @Nullable
+    private static List<Predicate> filtersWithin(
+            @Nullable List<Predicate> filters, Set<String> fields) {
+        if (isNullOrEmpty(filters)) {
+            return filters;
+        }
+
+        return filters.stream()
+                .filter(filter -> 
fields.containsAll(collectFieldNames(filter)))
+                .collect(Collectors.toList());
+    }
+
+    /**
+     * Devolve unconditionally, not only on a schema id mismatch: a data 
evolution file schema is
+     * projected to the columns the file actually wrote, so it is a subset of 
the table fields even
+     * for the same schema id. Filters on columns the file does not contain 
are dropped here.
+     */
+    private List<Predicate> devolveFilters(
+            @Nullable List<Predicate> filters, TableSchema dataSchema) {
+        List<Predicate> dataFilters =
+                SchemaEvolutionUtil.devolveFilters(
+                        schema.fields(), dataSchema.fields(), filters, false);
+        return excludePredicateWithFields(dataFilters, new 
HashSet<>(dataSchema.partitionKeys()));
+    }
+
     @Nullable
     private DeletionVectorWithRange readDeletionVector(
             List<DataFileMeta> group, @Nullable DeletionVector.Factory 
deletionVectorFactory)
@@ -680,6 +906,52 @@ public class DataEvolutionSplitRead implements 
SplitRead<InternalRow> {
         }
     }
 
+    /**
+     * Key of a single file reader mapping. Both the fields the mapping reads 
and the filters it
+     * pushes down depend on the read type, which is not the same for every 
split: {@link
+     * IndexedSplitRecordReader#readInfo} adds a row id to the read type when 
the split carries
+     * scores. The columns the file wrote are part of the key as well, they 
decide which filters the
+     * file can answer.
+     */
+    private static class SingleFileKey {
+
+        private final long schemaId;
+        private final String formatIdentifier;
+        @Nullable private final List<String> writeCols;
+        private final RowType readRowType;
+
+        private SingleFileKey(
+                long schemaId,
+                String formatIdentifier,
+                @Nullable List<String> writeCols,
+                RowType readRowType) {
+            this.schemaId = schemaId;
+            this.formatIdentifier = formatIdentifier;
+            this.writeCols = writeCols;
+            this.readRowType = readRowType;
+        }
+
+        @Override
+        public boolean equals(Object o) {
+            if (this == o) {
+                return true;
+            }
+            if (!(o instanceof SingleFileKey)) {
+                return false;
+            }
+            SingleFileKey that = (SingleFileKey) o;
+            return schemaId == that.schemaId
+                    && Objects.equals(formatIdentifier, that.formatIdentifier)
+                    && Objects.equals(writeCols, that.writeCols)
+                    && Objects.equals(readRowType, that.readRowType);
+        }
+
+        @Override
+        public int hashCode() {
+            return Objects.hash(schemaId, formatIdentifier, writeCols, 
readRowType);
+        }
+    }
+
     private static class DeletionVectorWithRange {
 
         private final Range range;
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionFileIndexTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionFileIndexTest.java
new file mode 100644
index 0000000000..178e5ed9e2
--- /dev/null
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionFileIndexTest.java
@@ -0,0 +1,805 @@
+/*
+ * 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.table;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.serializer.InternalRowSerializer;
+import org.apache.paimon.deletionvectors.BitmapDeletionVector;
+import org.apache.paimon.deletionvectors.DeletionVector;
+import org.apache.paimon.deletionvectors.append.BaseAppendDeleteFileMaintainer;
+import org.apache.paimon.fileindex.FileIndexOptions;
+import org.apache.paimon.fileindex.bitmap.BitmapFileIndexFactory;
+import org.apache.paimon.fileindex.bloomfilter.BloomFilterFileIndexFactory;
+import org.apache.paimon.index.IndexFileMeta;
+import org.apache.paimon.io.CompactIncrement;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.io.DataIncrement;
+import org.apache.paimon.manifest.FileKind;
+import org.apache.paimon.manifest.IndexManifestEntry;
+import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.reader.RecordReader;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaChange;
+import org.apache.paimon.table.sink.BatchTableCommit;
+import org.apache.paimon.table.sink.BatchTableWrite;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageImpl;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.InnerTableRead;
+import org.apache.paimon.table.source.ReadBuilder;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.table.source.TableRead;
+import org.apache.paimon.table.source.TableScan;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import javax.annotation.Nullable;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import static org.apache.paimon.table.SpecialFields.rowTypeWithRowId;
+import static org.apache.paimon.utils.DataEvolutionUtils.retrieveAnchorFile;
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * Tests file index and predicate push down in {@link
+ * org.apache.paimon.operation.DataEvolutionSplitRead}, on two levels.
+ *
+ * <p>{@link #readWithFilter} creates the readers directly from the splits, 
without {@link
+ * TableRead#executeFilter()}, so an empty result proves that the reader 
itself pruned the rows.
+ *
+ * <p>{@link #query} reads the whole plan through {@link 
TableRead#executeFilter()}, the residual
+ * filter both Flink and Spark keep on top of the push down, so it returns the 
answer a query gives.
+ * Those tests assert the rows themselves, and compare them with {@link 
#fullScanFiltered}: the push
+ * down is an optimization and must never change the result. A row count alone 
can not say that, it
+ * passes just as well when nothing is pushed down at all.
+ *
+ * <p>A filter on a column that is not part of the read type has no {@link 
#query} counterpart:
+ * {@link TableRead#executeFilter()} can not project such a predicate onto the 
read row and silently
+ * keeps every row, so only the split level assertion says anything.
+ *
+ * <p>Filter values are always picked inside the min/max range of the column, 
otherwise the group
+ * level stats pruning in {@link 
org.apache.paimon.operation.DataEvolutionFileStoreScan} would drop
+ * the split before the reader ever sees it.
+ */
+public class DataEvolutionFileIndexTest extends DataEvolutionTestBase {
+
+    private static final int ROW_COUNT = 100;
+
+    /** Inside the min/max range of f1 ("a000".."a099"), but not written. */
+    private static final String MISSING_F1 = "a050x";
+
+    /** Inside the min/max range of f2 ("b000".."b099"), but not written. */
+    private static final String MISSING_F2 = "b050x";
+
+    @Test
+    public void testSingleFileSkippedByFileIndex() throws Exception {
+        // a standalone .index file, only the reader can evaluate it
+        FileStoreTable table = createTable("single_file", bloomOptions("f1", 
"1 B"));
+        writeAllColumns(table, ROW_COUNT);
+
+        assertThat(readWithFilter(table, equalF1(MISSING_F1))).isEmpty();
+
+        List<InternalRow> hit = readWithFilter(table, equalF1(f1(50)));
+        assertThat(hit).isNotEmpty();
+        assertThat(hit).anyMatch(row -> row.getInt(0) == 50);
+        assertAligned(hit);
+    }
+
+    @Test
+    public void testQuerySingleFileWithFileIndex() throws Exception {
+        FileStoreTable table = createTable("single_file_query", 
bloomOptions("f1", "1 B"));
+        writeAllColumns(table, ROW_COUNT);
+
+        Predicate filter = equalF1(f1(50));
+        List<InternalRow> rows = query(table, filter);
+        // the push down is an optimization, it must not change the answer
+        
assertThat(rows).containsExactlyInAnyOrderElementsOf(fullScanFiltered(table, 
filter));
+        assertRow(assertSingleRow(rows), 50);
+
+        assertThat(query(table, equalF1(MISSING_F1))).isEmpty();
+    }
+
+    @Test
+    public void testMergedGroupSkippedByFileIndex() throws Exception {
+        // an embedded index, the default threshold keeps it in the manifest
+        FileStoreTable table = createTable("merged_group", 
Collections.emptyMap());
+        writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), 
bloomOptions("f2", null));
+        assertMergedGroup(table);
+
+        assertThat(readWithFilter(table, equalF2(MISSING_F2))).isEmpty();
+    }
+
+    @Test
+    public void testMergedGroupKeepsColumnsAligned() throws Exception {
+        FileStoreTable table = createTable("merged_aligned", 
Collections.emptyMap());
+        writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), 
bloomOptions("f2", null));
+        assertMergedGroup(table);
+
+        // a merged group is never row filtered, but the rows it returns must 
stay aligned
+        List<InternalRow> rows = readWithFilter(table, equalF2(f2(50)));
+        assertThat(rows).hasSize(ROW_COUNT);
+        assertAligned(rows);
+    }
+
+    @Test
+    public void testQueryMergedGroup() throws Exception {
+        FileStoreTable table = createTable("merged_group_query", 
Collections.emptyMap());
+        writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), 
bloomOptions("f2", null));
+        assertMergedGroup(table);
+
+        // the group is read whole, the residual filter has to cut it down to 
the matching row
+        Predicate filter = equalF2(f2(50));
+        List<InternalRow> rows = query(table, filter);
+        
assertThat(rows).containsExactlyInAnyOrderElementsOf(fullScanFiltered(table, 
filter));
+        assertRow(assertSingleRow(rows), 50);
+
+        assertThat(query(table, equalF2(MISSING_F2))).isEmpty();
+    }
+
+    @Test
+    public void testBitmapIndexSelectsMatchingRowsOnly() throws Exception {
+        FileStoreTable table = createTable("bitmap", bitmapOptions("f1"));
+        writeAllColumns(table, ROW_COUNT);
+
+        // a bitmap index is exact, the reader returns the matching row and 
nothing else
+        List<InternalRow> rows = readWithFilter(table, equalF1(f1(50)));
+        assertThat(rows).hasSize(1);
+        assertThat(rows.get(0).getInt(0)).isEqualTo(50);
+        assertAligned(rows);
+
+        assertThat(readWithFilter(table, equalF1(MISSING_F1))).isEmpty();
+    }
+
+    @Test
+    public void testQueryWithBitmapIndex() throws Exception {
+        FileStoreTable table = createTable("bitmap_query", 
bitmapOptions("f1"));
+        writeAllColumns(table, ROW_COUNT);
+
+        Predicate filter = equalF1(f1(50));
+        List<InternalRow> rows = query(table, filter);
+        
assertThat(rows).containsExactlyInAnyOrderElementsOf(fullScanFiltered(table, 
filter));
+        assertRow(assertSingleRow(rows), 50);
+
+        assertThat(query(table, equalF1(MISSING_F1))).isEmpty();
+    }
+
+    @Test
+    public void testMergedGroupSkippedAfterColumnRename() throws Exception {
+        FileStoreTable table = createTable("renamed", Collections.emptyMap());
+        writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), 
bloomOptions("f2", null));
+        catalog.alterTable(identifier("renamed"), 
SchemaChange.renameColumn("f2", "f3"), false);
+
+        // the index of the file was written under the old column name, the 
filter has to be
+        // devolved by field id before it can be evaluated
+        FileStoreTable renamed = getTable(identifier("renamed"));
+        PredicateBuilder builder = new PredicateBuilder(renamed.rowType());
+        int f3 = renamed.rowType().getFieldIndex("f3");
+
+        assertThat(readWithFilter(renamed, builder.equal(f3, 
BinaryString.fromString(MISSING_F2))))
+                .isEmpty();
+        assertThat(readWithFilter(renamed, builder.equal(f3, 
BinaryString.fromString(f2(50)))))
+                .hasSize(ROW_COUNT);
+    }
+
+    @Test
+    public void testQueryAfterColumnRename() throws Exception {
+        FileStoreTable table = createTable("renamed_query", 
Collections.emptyMap());
+        writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), 
bloomOptions("f2", null));
+        catalog.alterTable(
+                identifier("renamed_query"), SchemaChange.renameColumn("f2", 
"f3"), false);
+
+        FileStoreTable renamed = getTable(identifier("renamed_query"));
+        PredicateBuilder builder = new PredicateBuilder(renamed.rowType());
+        int f3 = renamed.rowType().getFieldIndex("f3");
+
+        Predicate filter = builder.equal(f3, BinaryString.fromString(f2(50)));
+        List<InternalRow> rows = query(renamed, filter);
+        
assertThat(rows).containsExactlyInAnyOrderElementsOf(fullScanFiltered(renamed, 
filter));
+        assertRow(assertSingleRow(rows), 50);
+
+        assertThat(query(renamed, builder.equal(f3, 
BinaryString.fromString(MISSING_F2))))
+                .isEmpty();
+    }
+
+    @Test
+    public void testFileIndexReadDisabled() throws Exception {
+        Map<String, String> options = bloomOptions("f1", "1 B");
+        options.put(CoreOptions.FILE_INDEX_READ_ENABLED.key(), "false");
+        FileStoreTable table = createTable("index_disabled", options);
+        writeAllColumns(table, ROW_COUNT);
+
+        assertThat(readWithFilter(table, 
equalF1(MISSING_F1))).hasSize(ROW_COUNT);
+    }
+
+    @Test
+    public void testQueryWithFileIndexReadDisabled() throws Exception {
+        Map<String, String> options = bloomOptions("f1", "1 B");
+        options.put(CoreOptions.FILE_INDEX_READ_ENABLED.key(), "false");
+        FileStoreTable table = createTable("index_disabled_query", options);
+        writeAllColumns(table, ROW_COUNT);
+
+        // the index only decides how much is read, never what the query 
answers
+        Predicate filter = equalF1(f1(50));
+        List<InternalRow> rows = query(table, filter);
+        
assertThat(rows).containsExactlyInAnyOrderElementsOf(fullScanFiltered(table, 
filter));
+        assertRow(assertSingleRow(rows), 50);
+
+        assertThat(query(table, equalF1(MISSING_F1))).isEmpty();
+    }
+
+    @Test
+    public void testRowTrackingFilterIsNotPushedDown() throws Exception {
+        FileStoreTable table = createTable("row_tracking", 
Collections.emptyMap());
+        writeSplitColumns(table, ROW_COUNT, bloomOptions("f1", "1 B"), 
Collections.emptyMap());
+
+        // _ROW_ID is assigned from the manifest entry, a filter on it must 
not reach the file
+        // index or the format reader, and must not break the filter 
devolution either, which
+        // resolves predicate fields against the table schema and knows 
nothing about system fields
+        PredicateBuilder builder = new 
PredicateBuilder(rowTypeWithRowId(rowType()));
+        Predicate predicate =
+                PredicateBuilder.and(
+                        builder.between(3, 0L, (long) ROW_COUNT),
+                        builder.equal(1, BinaryString.fromString(f1(50))));
+
+        List<InternalRow> rows = readWithFilter(table, predicate);
+        assertThat(rows).hasSize(ROW_COUNT);
+        assertAligned(rows);
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"parquet", "orc", "avro"})
+    public void testFilterOnColumnMissingFromFileIsNotApplied(String format) 
throws Exception {
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.FILE_FORMAT.key(), format);
+        options.put(CoreOptions.DELETION_VECTORS_ENABLED.key(), "false");
+        FileStoreTable table = createTable("missing_column_" + format, 
options);
+        writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), 
Collections.emptyMap());
+
+        // projecting f2 leaves only the file that does not contain f1, so the 
filter reaches a
+        // format reader that knows nothing about the column and must not drop 
anything
+        List<InternalRow> rows = readWithFilter(table, equalF1(f1(50)), 
rowType().project("f2"));
+        assertThat(rows).hasSize(ROW_COUNT);
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"parquet", "orc", "avro"})
+    public void testReadOnlyRowIdWithFilterOnDataColumn(String format) throws 
Exception {
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.FILE_FORMAT.key(), format);
+        FileStoreTable table = createTable("row_id_only_" + format, options);
+        writeAllColumns(table, ROW_COUNT);
+
+        // the read type holds no data column at all, so the filter can not be 
pushed to the
+        // format: a format only evaluates a predicate over the columns it was 
asked to read and
+        // treats the others as absent, which drops every row instead of none
+        List<InternalRow> rows =
+                readWithFilter(table, equalF1(f1(50)), 
RowType.of(SpecialFields.ROW_ID));
+        assertThat(rowIds(rows)).containsExactlyInAnyOrderElementsOf(rowIds(0, 
ROW_COUNT));
+    }
+
+    @Test
+    public void testProjectionWithoutFilterColumn() throws Exception {
+        FileStoreTable table = createTable("projection_without_filter", 
Collections.emptyMap());
+        writeAllColumns(table, ROW_COUNT);
+
+        // f1 is filtered but not read, nothing proves anything about the 
rows, all of them stay
+        List<InternalRow> rows = readWithFilter(table, equalF1(f1(50)), 
rowType().project("f0"));
+        assertThat(rows).hasSize(ROW_COUNT);
+        assertThat(rows).anyMatch(row -> row.getInt(0) == 50);
+    }
+
+    @Test
+    public void testFileIndexIsGivenUpForAColumnOutsideTheReadType() throws 
Exception {
+        FileStoreTable table = createTable("projection_index", 
bloomOptions("f1", "1 B"));
+        writeAllColumns(table, ROW_COUNT);
+
+        // the index of this file does prove that no row of it matches, but a 
file only owns the
+        // columns of the read type: the file that owns f1 can be pruned out 
of the split, see
+        // testProjectionPruningAwayTheWinnerOfTheFilterColumn, so the push 
down gives the column
+        // up rather than trust whichever file is left holding it
+        assertThat(readWithFilter(table, equalF1(MISSING_F1), 
rowType().project("f0")))
+                .hasSize(ROW_COUNT);
+
+        // reading f1 puts it back in reach
+        assertThat(readWithFilter(table, equalF1(MISSING_F1), 
rowType().project("f0", "f1")))
+                .isEmpty();
+    }
+
+    @Test
+    public void testBitmapSelectionReturnsExactRowId() throws Exception {
+        FileStoreTable table = createTable("bitmap_row_id", 
bitmapOptions("f1"));
+        writeAllColumns(table, ROW_COUNT);
+        // a second row id range, so a row id can not be confused with a 
position inside a file
+        writeAllColumns(table, ROW_COUNT);
+
+        // the bitmap is exact, and the row id of what it selects comes from 
the position the row
+        // has in its file, not from its position in the batch the selection 
left over
+        RowType readType = 
rowTypeWithRowId(rowType()).project(SpecialFields.ROW_ID.name(), "f1");
+        List<InternalRow> rows = readWithFilter(table, equalF1(f1(50)), 
readType);
+        assertThat(rowIds(rows)).containsExactlyInAnyOrder(50L, ROW_COUNT + 
50L);
+    }
+
+    @Test
+    public void testProjectionPruningAwayTheWinnerOfTheFilterColumn() throws 
Exception {
+        FileStoreTable table = createTable("pruned_winner", 
Collections.emptyMap());
+        writeThenOverwriteF1(table, ROW_COUNT);
+
+        // f1 was rewritten as c* by a second file, and projecting f0 prunes 
that file out of the
+        // split, see DataEvolutionFileStoreScan#pruneByReadType. The old file 
is left holding a*
+        // and a bloom index that knows nothing about c*, so nothing about it 
may be used to prove
+        // that a row does not match: the row does match, through the file 
that is not there.
+        List<InternalRow> rows = readWithFilter(table, equalF1(c1(50)), 
rowType().project("f0"));
+        assertThat(rows).anyMatch(row -> row.getInt(0) == 50);
+    }
+
+    @Test
+    public void testEmptyReadTypeWithFilter() throws Exception {
+        FileStoreTable table = createTable("empty_read_type", 
Collections.emptyMap());
+        writeAllColumns(table, ROW_COUNT);
+
+        // a count(*) read projects nothing, so nothing can be pushed to the 
format either
+        assertThat(countWithFilter(table, equalF1(f1(50)), 
RowType.of())).isEqualTo(ROW_COUNT);
+    }
+
+    @Test
+    public void testReaderMappingIsNotSharedBetweenReadTypes() throws 
Exception {
+        FileStoreTable table = createTable("read_type_reuse", 
Collections.emptyMap());
+        writeAllColumns(table, ROW_COUNT);
+
+        FileStoreTable latest = getTable(identifier(table.name()));
+        TableScan.Plan plan = latest.newReadBuilder().newScan().plan();
+        InnerTableRead read = latest.newRead();
+
+        // a read can be reconfigured after it created readers, see 
AppendTableRead#applyReadType,
+        // so a reader mapping built for one read type must not answer for 
another
+        RowType first = rowType().project("f1");
+        read.withReadType(first);
+        assertThat(collect(read, plan, 
first).get(0).getString(0).toString()).isEqualTo(f1(0));
+
+        RowType second = rowType().project("f2");
+        read.withReadType(second);
+        assertThat(collect(read, plan, 
second).get(0).getString(0).toString()).isEqualTo(f2(0));
+    }
+
+    @Test
+    public void testReaderMappingIsNotSharedBetweenFilters() throws Exception {
+        FileStoreTable table = createTable("filter_reuse", 
Collections.emptyMap());
+        writeAllColumns(table, ROW_COUNT);
+
+        FileStoreTable latest = getTable(identifier(table.name()));
+        // planned without a filter, so only the readers can drop anything here
+        TableScan.Plan plan = latest.newReadBuilder().newScan().plan();
+        InnerTableRead read = latest.newRead();
+
+        // outside the min/max of f1, the format reader prunes the whole file
+        read.withFilter(equalF1("z999"));
+        assertThat(collect(read, plan, latest.rowType())).isEmpty();
+
+        // the mapping above carries that filter, the second read must not 
inherit it
+        read.withFilter(equalF1(f1(50)));
+        assertThat(collect(read, plan, latest.rowType())).isNotEmpty();
+    }
+
+    @Test
+    public void testBitmapSelectionComposesWithDeletionVector() throws 
Exception {
+        Map<String, String> options = bitmapOptions("f1");
+        options.put(CoreOptions.DELETION_VECTORS_ENABLED.key(), "true");
+        FileStoreTable table = createTable("bitmap_dv", options);
+        writeAllColumns(table, ROW_COUNT);
+
+        // both the bitmap selection and the deletion vector filter by 
position, so deleting the
+        // very row the bitmap selects has to make it disappear
+        deleteRows(table, 50);
+        assertThat(readWithFilter(table, equalF1(f1(50)))).isEmpty();
+
+        // deleting a neighbour must leave the selected row untouched
+        FileStoreTable other = createTable("bitmap_dv_other", options);
+        writeAllColumns(other, ROW_COUNT);
+        deleteRows(other, 51);
+        List<InternalRow> rows = readWithFilter(other, equalF1(f1(50)));
+        assertThat(rows).hasSize(1);
+        assertThat(rows.get(0).getInt(0)).isEqualTo(50);
+    }
+
+    @Test
+    public void testQueryComposesBitmapWithDeletionVector() throws Exception {
+        Map<String, String> options = bitmapOptions("f1");
+        options.put(CoreOptions.DELETION_VECTORS_ENABLED.key(), "true");
+        FileStoreTable table = createTable("bitmap_dv_query", options);
+        writeAllColumns(table, ROW_COUNT);
+        deleteRows(table, 50);
+
+        assertThat(query(table, equalF1(f1(50)))).isEmpty();
+
+        FileStoreTable other = createTable("bitmap_dv_other_query", options);
+        writeAllColumns(other, ROW_COUNT);
+        deleteRows(other, 51);
+
+        Predicate filter = equalF1(f1(50));
+        List<InternalRow> rows = query(other, filter);
+        
assertThat(rows).containsExactlyInAnyOrderElementsOf(fullScanFiltered(other, 
filter));
+        assertRow(assertSingleRow(rows), 50);
+    }
+
+    @Test
+    public void testMergedGroupKeptWhenFilterColumnOverwritten() throws 
Exception {
+        FileStoreTable table = createTable("overwritten", 
Collections.emptyMap());
+        writeThenOverwriteF1(table, ROW_COUNT);
+
+        // f1 was rewritten by a newer file (c*); the old file still holds a* 
and a bloom index that
+        // knows nothing about c*, so its index must not veto the group
+        Predicate filter = equalF1(c1(50));
+        List<InternalRow> rows = query(table, filter);
+        
assertThat(rows).containsExactlyInAnyOrderElementsOf(fullScanFiltered(table, 
filter));
+
+        InternalRow row = assertSingleRow(rows);
+        assertThat(row.getInt(0)).isEqualTo(50);
+        assertThat(row.getString(1).toString()).isEqualTo(c1(50));
+        // no file of the group wrote f2
+        assertThat(row.isNullAt(2)).isTrue();
+
+        // and the overwritten value must be gone, the column merge has to 
take f1 from the new file
+        assertThat(query(table, equalF1(f1(50)))).isEmpty();
+    }
+
+    /** Commits a deletion vector for the anchor file of the only row id group 
of {@code table}. */
+    private void deleteRows(FileStoreTable table, long... positions) throws 
Exception {
+        FileStoreTable latest = getTable(identifier(table.name()));
+        List<DataFileMeta> dataFiles =
+                ((DataSplit) 
latest.newReadBuilder().newScan().plan().splits().get(0)).dataFiles();
+        String anchor = retrieveAnchorFile(dataFiles, file -> file).fileName();
+
+        BaseAppendDeleteFileMaintainer maintainer =
+                BaseAppendDeleteFileMaintainer.forUnawareAppend(
+                        latest.store().newIndexFileHandler(),
+                        latest.latestSnapshot().get(),
+                        BinaryRow.EMPTY_ROW);
+        DeletionVector deletionVector = new BitmapDeletionVector();
+        for (long position : positions) {
+            deletionVector.delete(position);
+        }
+        maintainer.notifyNewDeletionVector(anchor, deletionVector);
+
+        List<IndexFileMeta> newIndexFiles = new ArrayList<>();
+        for (IndexManifestEntry entry : maintainer.persist()) {
+            if (entry.kind() == FileKind.ADD) {
+                newIndexFiles.add(entry.indexFile());
+            }
+        }
+
+        try (BatchTableCommit commit = 
latest.newBatchWriteBuilder().newCommit()) {
+            commit.commit(
+                    Collections.singletonList(
+                            new CommitMessageImpl(
+                                    BinaryRow.EMPTY_ROW,
+                                    BucketMode.UNAWARE_BUCKET,
+                                    null,
+                                    new DataIncrement(
+                                            Collections.emptyList(),
+                                            Collections.emptyList(),
+                                            Collections.emptyList(),
+                                            newIndexFiles,
+                                            Collections.emptyList()),
+                                    CompactIncrement.emptyIncrement())));
+        }
+    }
+
+    private FileStoreTable createTable(String name, Map<String, String> 
options) throws Exception {
+        Schema.Builder builder =
+                Schema.newBuilder()
+                        .column("f0", DataTypes.INT())
+                        .column("f1", DataTypes.STRING())
+                        .column("f2", DataTypes.STRING())
+                        .option(CoreOptions.ROW_TRACKING_ENABLED.key(), "true")
+                        .option(CoreOptions.DATA_EVOLUTION_ENABLED.key(), 
"true");
+        options.forEach(builder::option);
+        Identifier identifier = identifier(name);
+        catalog.createTable(identifier, builder.build(), false);
+        return getTable(identifier);
+    }
+
+    private static Map<String, String> bloomOptions(String column, String 
inManifestThreshold) {
+        String prefix =
+                FileIndexOptions.FILE_INDEX + "." + 
BloomFilterFileIndexFactory.BLOOM_FILTER + ".";
+        Map<String, String> options = new HashMap<>();
+        options.put(prefix + CoreOptions.COLUMNS, column);
+        options.put(prefix + column + ".items", "200");
+        options.put(prefix + column + ".fpp", "0.001");
+        if (inManifestThreshold != null) {
+            options.put(CoreOptions.FILE_INDEX_IN_MANIFEST_THRESHOLD.key(), 
inManifestThreshold);
+        }
+        return options;
+    }
+
+    private static Map<String, String> bitmapOptions(String column) {
+        Map<String, String> options = new HashMap<>();
+        options.put(
+                FileIndexOptions.FILE_INDEX
+                        + "."
+                        + BitmapFileIndexFactory.BITMAP_INDEX
+                        + "."
+                        + CoreOptions.COLUMNS,
+                column);
+        return options;
+    }
+
+    private void writeAllColumns(FileStoreTable table, int count) throws 
Exception {
+        BatchWriteBuilder builder = table.newBatchWriteBuilder();
+        try (BatchTableWrite write = builder.newWrite();
+                BatchTableCommit commit = builder.newCommit()) {
+            for (int i = 0; i < count; i++) {
+                write.write(
+                        GenericRow.of(
+                                i, BinaryString.fromString(f1(i)), 
BinaryString.fromString(f2(i))));
+            }
+            commit.commit(write.prepareCommit());
+        }
+    }
+
+    /**
+     * Writes f0/f1 and f2 into two files sharing one row id range, so they 
have to be merged. The
+     * index options are per write, a file index can only be configured for 
columns the write
+     * actually contains (see {@link 
org.apache.paimon.io.DataFileIndexWriter}).
+     */
+    private void writeSplitColumns(
+            FileStoreTable table,
+            int count,
+            Map<String, String> firstOptions,
+            Map<String, String> secondOptions)
+            throws Exception {
+        RowType writeType0 = table.rowType().project(Arrays.asList("f0", 
"f1"));
+        RowType writeType1 = 
table.rowType().project(Collections.singletonList("f2"));
+
+        BatchWriteBuilder builder = 
table.copy(firstOptions).newBatchWriteBuilder();
+        try (BatchTableWrite write = 
builder.newWrite().withWriteType(writeType0);
+                BatchTableCommit commit = builder.newCommit()) {
+            for (int i = 0; i < count; i++) {
+                write.write(GenericRow.of(i, BinaryString.fromString(f1(i))));
+            }
+            commit.commit(write.prepareCommit());
+        }
+
+        FileStoreTable latest = getTable(identifier(table.name()));
+        long firstRowId = 
latest.snapshotManager().latestSnapshot().nextRowId() - count;
+        builder = latest.copy(secondOptions).newBatchWriteBuilder();
+        try (BatchTableWrite write = 
builder.newWrite().withWriteType(writeType1);
+                BatchTableCommit commit = builder.newCommit()) {
+            for (int i = 0; i < count; i++) {
+                write.write(GenericRow.of(BinaryString.fromString(f2(i))));
+            }
+            List<CommitMessage> commitables = write.prepareCommit();
+            setSingleFileFirstRowId(commitables, firstRowId);
+            commit.commit(commitables);
+        }
+    }
+
+    /**
+     * Writes f0/f1 (a*) with a bloom index on f1, then overwrites f1 with c* 
in a second file that
+     * shares the same row id range and has a higher sequence number, so the 
column merge takes f1
+     * from the newer file.
+     */
+    private void writeThenOverwriteF1(FileStoreTable table, int count) throws 
Exception {
+        RowType writeType0 = table.rowType().project(Arrays.asList("f0", 
"f1"));
+        BatchWriteBuilder builder = table.copy(bloomOptions("f1", 
null)).newBatchWriteBuilder();
+        try (BatchTableWrite write = 
builder.newWrite().withWriteType(writeType0);
+                BatchTableCommit commit = builder.newCommit()) {
+            for (int i = 0; i < count; i++) {
+                write.write(GenericRow.of(i, BinaryString.fromString(f1(i))));
+            }
+            commit.commit(write.prepareCommit());
+        }
+
+        FileStoreTable latest = getTable(identifier(table.name()));
+        long firstRowId = 
latest.snapshotManager().latestSnapshot().nextRowId() - count;
+        RowType writeType1 = 
table.rowType().project(Collections.singletonList("f1"));
+        builder = latest.newBatchWriteBuilder();
+        try (BatchTableWrite write = 
builder.newWrite().withWriteType(writeType1);
+                BatchTableCommit commit = builder.newCommit()) {
+            for (int i = 0; i < count; i++) {
+                write.write(GenericRow.of(BinaryString.fromString(c1(i))));
+            }
+            List<CommitMessage> commitables = write.prepareCommit();
+            setSingleFileFirstRowId(commitables, firstRowId);
+            commit.commit(commitables);
+        }
+    }
+
+    /**
+     * {@link #setFirstRowId} stamps the same first row id on every file of 
the commit, which only
+     * describes a valid layout when the write produced a single file.
+     */
+    private void setSingleFileFirstRowId(List<CommitMessage> commitables, long 
firstRowId) {
+        List<DataFileMeta> files = new ArrayList<>();
+        for (CommitMessage commitable : commitables) {
+            files.addAll(((CommitMessageImpl) 
commitable).newFilesIncrement().newFiles());
+        }
+        assertThat(files).hasSize(1);
+        setFirstRowId(commitables, firstRowId);
+    }
+
+    private List<InternalRow> readWithFilter(FileStoreTable table, Predicate 
predicate)
+            throws Exception {
+        return readWithFilter(table, predicate, null);
+    }
+
+    private List<InternalRow> readWithFilter(
+            FileStoreTable table, Predicate predicate, @Nullable RowType 
readType)
+            throws Exception {
+        FileStoreTable latest = getTable(identifier(table.name()));
+        ReadBuilder readBuilder = 
latest.newReadBuilder().withFilter(predicate);
+        if (readType != null) {
+            readBuilder.withReadType(readType);
+        }
+        TableRead read = readBuilder.newRead();
+        InternalRowSerializer serializer =
+                new InternalRowSerializer(readType == null ? latest.rowType() 
: readType);
+        List<InternalRow> rows = new ArrayList<>();
+        for (Split split : readBuilder.newScan().plan().splits()) {
+            try (RecordReader<InternalRow> reader = read.createReader(split)) {
+                reader.forEachRemaining(row -> rows.add(serializer.copy(row)));
+            }
+        }
+        return rows;
+    }
+
+    /** Like {@link #readWithFilter}, but only counts, so an empty read type 
can be used. */
+    private int countWithFilter(FileStoreTable table, Predicate predicate, 
RowType readType)
+            throws Exception {
+        FileStoreTable latest = getTable(identifier(table.name()));
+        ReadBuilder readBuilder =
+                
latest.newReadBuilder().withFilter(predicate).withReadType(readType);
+        TableRead read = readBuilder.newRead();
+        int count = 0;
+        for (Split split : readBuilder.newScan().plan().splits()) {
+            try (RecordReader<InternalRow> reader = read.createReader(split)) {
+                RecordReader.RecordIterator<InternalRow> batch;
+                while ((batch = reader.readBatch()) != null) {
+                    while (batch.next() != null) {
+                        count++;
+                    }
+                    batch.releaseBatch();
+                }
+            }
+        }
+        return count;
+    }
+
+    /**
+     * Reads the whole plan the way an engine does, with the residual filter 
on top of the push
+     * down.
+     */
+    private List<InternalRow> query(FileStoreTable table, Predicate predicate) 
throws Exception {
+        FileStoreTable latest = getTable(identifier(table.name()));
+        ReadBuilder readBuilder = 
latest.newReadBuilder().withFilter(predicate);
+        return collect(
+                readBuilder.newRead().executeFilter(),
+                readBuilder.newScan().plan(),
+                latest.rowType());
+    }
+
+    /** What {@link #query} has to return: a full scan filtered in memory, 
with no push down. */
+    private List<InternalRow> fullScanFiltered(FileStoreTable table, Predicate 
predicate)
+            throws Exception {
+        FileStoreTable latest = getTable(identifier(table.name()));
+        ReadBuilder readBuilder = latest.newReadBuilder();
+        List<InternalRow> rows = new ArrayList<>();
+        for (InternalRow row :
+                collect(readBuilder.newRead(), readBuilder.newScan().plan(), 
latest.rowType())) {
+            if (predicate.test(row)) {
+                rows.add(row);
+            }
+        }
+        return rows;
+    }
+
+    /**
+     * Materializes as {@link BinaryRow}, the reader hands out row 
implementations without a value
+     * based equals.
+     */
+    private static List<InternalRow> collect(TableRead read, TableScan.Plan 
plan, RowType rowType)
+            throws Exception {
+        InternalRowSerializer serializer = new InternalRowSerializer(rowType);
+        List<InternalRow> rows = new ArrayList<>();
+        try (RecordReader<InternalRow> reader = read.createReader(plan)) {
+            reader.forEachRemaining(row -> 
rows.add(serializer.toBinaryRow(row).copy()));
+        }
+        return rows;
+    }
+
+    private void assertMergedGroup(FileStoreTable table) throws Exception {
+        FileStoreTable latest = getTable(identifier(table.name()));
+        List<Split> splits = latest.newReadBuilder().newScan().plan().splits();
+        assertThat(splits).hasSize(1);
+        assertThat(((DataSplit) splits.get(0)).dataFiles()).hasSize(2);
+    }
+
+    /** Row ids of rows read with {@code RowType.of(SpecialFields.ROW_ID)}. */
+    private static List<Long> rowIds(List<InternalRow> rows) {
+        return rows.stream().map(row -> 
row.getLong(0)).collect(Collectors.toList());
+    }
+
+    private static List<Long> rowIds(long from, int count) {
+        List<Long> ids = new ArrayList<>(count);
+        for (int i = 0; i < count; i++) {
+            ids.add(from + i);
+        }
+        return ids;
+    }
+
+    private static InternalRow assertSingleRow(List<InternalRow> rows) {
+        assertThat(rows).hasSize(1);
+        return rows.get(0);
+    }
+
+    private static void assertRow(InternalRow row, int i) {
+        assertThat(row.getInt(0)).isEqualTo(i);
+        assertThat(row.getString(1).toString()).isEqualTo(f1(i));
+        assertThat(row.getString(2).toString()).isEqualTo(f2(i));
+    }
+
+    private static void assertAligned(List<InternalRow> rows) {
+        for (InternalRow row : rows) {
+            int f0 = row.getInt(0);
+            assertThat(row.getString(1).toString()).isEqualTo(f1(f0));
+            assertThat(row.getString(2).toString()).isEqualTo(f2(f0));
+        }
+    }
+
+    private static Predicate equalF1(String value) {
+        return new PredicateBuilder(rowType()).equal(1, 
BinaryString.fromString(value));
+    }
+
+    private static Predicate equalF2(String value) {
+        return new PredicateBuilder(rowType()).equal(2, 
BinaryString.fromString(value));
+    }
+
+    private static RowType rowType() {
+        return RowType.of(DataTypes.INT(), DataTypes.STRING(), 
DataTypes.STRING());
+    }
+
+    private static String f1(int i) {
+        return String.format("a%03d", i);
+    }
+
+    private static String f2(int i) {
+        return String.format("b%03d", i);
+    }
+
+    private static String c1(int i) {
+        return String.format("c%03d", i);
+    }
+}

Reply via email to