This is an automated email from the ASF dual-hosted git repository.

zhouyuan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git


The following commit(s) were added to refs/heads/main by this push:
     new 51a3c4e77a [VL] feat(iceberg): Add delete file bounds propagation 
(#12420)
51a3c4e77a is described below

commit 51a3c4e77a71ac2d6a3acf575e73be40cb0c4a3a
Author: inf <[email protected]>
AuthorDate: Fri Jul 31 08:08:51 2026 +0000

    [VL] feat(iceberg): Add delete file bounds propagation (#12420)
    
    Improve performance by skipping unrelated positional delete files when 
reading
---
 cpp/velox/compute/iceberg/IcebergPlanConverter.cc  |  27 ++-
 gluten-iceberg/pom.xml                             |   5 +
 .../substrait/rel/IcebergLocalFilesNode.java       |  39 +++-
 .../rel/IcebergLocalFilesNodeBoundsTest.java       | 214 +++++++++++++++++++++
 4 files changed, 280 insertions(+), 5 deletions(-)

diff --git a/cpp/velox/compute/iceberg/IcebergPlanConverter.cc 
b/cpp/velox/compute/iceberg/IcebergPlanConverter.cc
index 07c40e6e1c..5212381623 100644
--- a/cpp/velox/compute/iceberg/IcebergPlanConverter.cc
+++ b/cpp/velox/compute/iceberg/IcebergPlanConverter.cc
@@ -19,6 +19,24 @@
 
 namespace gluten {
 
+namespace {
+
+using SubstraitDeleteBoundsMap = 
::substrait::ReadRel_LocalFiles_FileOrFiles::IcebergReadOptions::DeleteFile::Map;
+
+std::unordered_map<int32_t, std::string> parseBounds(const 
SubstraitDeleteBoundsMap& bounds) {
+  std::unordered_map<int32_t, std::string> parsed;
+  parsed.reserve(bounds.key_values_size());
+
+  for (int i = 0; i < bounds.key_values_size(); ++i) {
+    const auto& kv = bounds.key_values(i);
+    parsed.emplace(kv.key(), kv.value());
+  }
+
+  return parsed;
+}
+
+} // namespace
+
 std::shared_ptr<IcebergSplitInfo> IcebergPlanConverter::parseIcebergSplitInfo(
     substrait::ReadRel_LocalFiles_FileOrFiles file,
     std::shared_ptr<SplitInfo> splitInfo) {
@@ -70,7 +88,14 @@ std::shared_ptr<IcebergSplitInfo> 
IcebergPlanConverter::parseIcebergSplitInfo(
           break;
       }
       deletes.emplace_back(IcebergDeleteFile(
-          fileContent, deleteFile.filepath(), format, 
deleteFile.recordcount(), deleteFile.filesize()));
+          fileContent,
+          deleteFile.filepath(),
+          format,
+          deleteFile.recordcount(),
+          deleteFile.filesize(),
+          {},
+          parseBounds(deleteFile.lowerbounds()),
+          parseBounds(deleteFile.upperbounds())));
     }
     icebergSplitInfo->deleteFilesVec.emplace_back(deletes);
   } else {
diff --git a/gluten-iceberg/pom.xml b/gluten-iceberg/pom.xml
index 056535b940..2cb694f8e9 100644
--- a/gluten-iceberg/pom.xml
+++ b/gluten-iceberg/pom.xml
@@ -101,6 +101,11 @@
       <artifactId>scalatest_${scala.binary.version}</artifactId>
       <scope>test</scope>
     </dependency>
+    <dependency>
+      <groupId>junit</groupId>
+      <artifactId>junit</artifactId>
+      <scope>test</scope>
+    </dependency>
   </dependencies>
 
   <build>
diff --git 
a/gluten-iceberg/src/main/java/org/apache/gluten/substrait/rel/IcebergLocalFilesNode.java
 
b/gluten-iceberg/src/main/java/org/apache/gluten/substrait/rel/IcebergLocalFilesNode.java
index cca3a4f79d..ac5841226a 100644
--- 
a/gluten-iceberg/src/main/java/org/apache/gluten/substrait/rel/IcebergLocalFilesNode.java
+++ 
b/gluten-iceberg/src/main/java/org/apache/gluten/substrait/rel/IcebergLocalFilesNode.java
@@ -18,11 +18,10 @@ package org.apache.gluten.substrait.rel;
 
 import io.substrait.proto.ReadRel;
 import org.apache.iceberg.DeleteFile;
+import org.apache.iceberg.FileContent;
 
-import java.util.ArrayList;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
+import java.nio.ByteBuffer;
+import java.util.*;
 
 public class IcebergLocalFilesNode extends LocalFilesNode {
   private final List<List<DeleteFile>> deleteFilesList;
@@ -96,6 +95,10 @@ public class IcebergLocalFilesNode extends LocalFilesNode {
       deleteFileBuilder.setFilePath(delete.path().toString());
       deleteFileBuilder.setFileSize(delete.fileSizeInBytes());
       deleteFileBuilder.setRecordCount(delete.recordCount());
+      if (delete.content() == FileContent.POSITION_DELETES) {
+        deleteFileBuilder.setLowerBounds(encodeBounds(delete.lowerBounds()));
+        deleteFileBuilder.setUpperBounds(encodeBounds(delete.upperBounds()));
+      }
       switch (delete.format()) {
         case PARQUET:
           ReadRel.LocalFiles.FileOrFiles.ParquetReadOptions parquetReadOptions 
=
@@ -118,4 +121,32 @@ public class IcebergLocalFilesNode extends LocalFilesNode {
     }
     fileBuilder.setIceberg(icebergBuilder);
   }
+
+  private static 
ReadRel.LocalFiles.FileOrFiles.IcebergReadOptions.DeleteFile.Map encodeBounds(
+      Map<Integer, ByteBuffer> bounds) {
+    ReadRel.LocalFiles.FileOrFiles.IcebergReadOptions.DeleteFile.Map.Builder 
builder =
+        
ReadRel.LocalFiles.FileOrFiles.IcebergReadOptions.DeleteFile.Map.newBuilder();
+
+    if (bounds == null || bounds.isEmpty()) {
+      return builder.build();
+    }
+
+    for (Map.Entry<Integer, ByteBuffer> entry : bounds.entrySet()) {
+      if (entry.getValue() == null) {
+        continue;
+      }
+
+      ByteBuffer duplicate = entry.getValue().asReadOnlyBuffer();
+      byte[] bytes = new byte[duplicate.remaining()];
+      duplicate.get(bytes);
+
+      builder.addKeyValues(
+          
ReadRel.LocalFiles.FileOrFiles.IcebergReadOptions.DeleteFile.Map.KeyValue.newBuilder()
+              .setKey(entry.getKey())
+              .setValue(Base64.getEncoder().encodeToString(bytes))
+              .build());
+    }
+
+    return builder.build();
+  }
 }
diff --git 
a/gluten-iceberg/src/test/java/org/apache/gluten/substrait/rel/IcebergLocalFilesNodeBoundsTest.java
 
b/gluten-iceberg/src/test/java/org/apache/gluten/substrait/rel/IcebergLocalFilesNodeBoundsTest.java
new file mode 100644
index 0000000000..01fc50035c
--- /dev/null
+++ 
b/gluten-iceberg/src/test/java/org/apache/gluten/substrait/rel/IcebergLocalFilesNodeBoundsTest.java
@@ -0,0 +1,214 @@
+/*
+ * 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.gluten.substrait.rel;
+
+import io.substrait.proto.ReadRel;
+import org.apache.iceberg.DeleteFile;
+import org.apache.iceberg.FileContent;
+import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.StructLike;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.nio.ByteBuffer;
+import java.util.Base64;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+public class IcebergLocalFilesNodeBoundsTest {
+  private static final int ICEBERG_POS_FIELD_ID = 2147483545;
+
+  @Test
+  public void serializesDeleteFileLowerAndUpperBounds() {
+    byte[] lowerBytes = new byte[] {1, 2, 3, 4};
+    byte[] upperBytes = new byte[] {5, 6, 7, 8};
+
+    FakeDeleteFile deleteFile =
+        new FakeDeleteFile(
+            FileContent.POSITION_DELETES,
+            FileFormat.PARQUET,
+            "/tmp/delete.parquet",
+            2L,
+            123L,
+            Collections.emptyList(),
+            Collections.singletonMap(ICEBERG_POS_FIELD_ID, 
ByteBuffer.wrap(lowerBytes)),
+            Collections.singletonMap(ICEBERG_POS_FIELD_ID, 
ByteBuffer.wrap(upperBytes)),
+            null);
+    IcebergLocalFilesNode node =
+        new IcebergLocalFilesNode(
+            0,
+            Collections.singletonList("/tmp/data.parquet"),
+            Collections.singletonList(0L),
+            Collections.singletonList(100L),
+            Collections.singletonList(Collections.emptyMap()),
+            LocalFilesNode.ReadFileFormat.ParquetReadFormat,
+            Collections.emptyList(),
+            Collections.singletonList(Collections.singletonList(deleteFile)),
+            Collections.singletonList(Collections.emptyMap()));
+
+    ReadRel.LocalFiles.FileOrFiles.Builder fileBuilder =
+        ReadRel.LocalFiles.FileOrFiles.newBuilder();
+
+    node.processFileBuilder(fileBuilder, 0);
+
+    ReadRel.LocalFiles.FileOrFiles.IcebergReadOptions.DeleteFile actualDelete =
+        fileBuilder.getIceberg().getDeleteFiles(0);
+
+    Assert.assertEquals(1, actualDelete.getLowerBounds().getKeyValuesCount());
+    Assert.assertEquals(
+        ICEBERG_POS_FIELD_ID, 
actualDelete.getLowerBounds().getKeyValues(0).getKey());
+    Assert.assertEquals(
+        Base64.getEncoder().encodeToString(lowerBytes),
+        actualDelete.getLowerBounds().getKeyValues(0).getValue());
+
+    Assert.assertEquals(1, actualDelete.getUpperBounds().getKeyValuesCount());
+    Assert.assertEquals(
+        ICEBERG_POS_FIELD_ID, 
actualDelete.getUpperBounds().getKeyValues(0).getKey());
+    Assert.assertEquals(
+        Base64.getEncoder().encodeToString(upperBytes),
+        actualDelete.getUpperBounds().getKeyValues(0).getValue());
+  }
+
+  private static final class FakeDeleteFile implements DeleteFile {
+    private final FileContent content;
+    private final FileFormat format;
+    private final String path;
+    private final long recordCount;
+    private final long fileSizeInBytes;
+    private final List<Integer> equalityFieldIds;
+    private final Map<Integer, ByteBuffer> lowerBounds;
+    private final Map<Integer, ByteBuffer> upperBounds;
+    private final Long dataSequenceNumber;
+
+    private FakeDeleteFile(
+        FileContent content,
+        FileFormat format,
+        String path,
+        long recordCount,
+        long fileSizeInBytes,
+        List<Integer> equalityFieldIds,
+        Map<Integer, ByteBuffer> lowerBounds,
+        Map<Integer, ByteBuffer> upperBounds,
+        Long dataSequenceNumber) {
+      this.content = content;
+      this.format = format;
+      this.path = path;
+      this.recordCount = recordCount;
+      this.fileSizeInBytes = fileSizeInBytes;
+      this.equalityFieldIds = equalityFieldIds;
+      this.lowerBounds = lowerBounds;
+      this.upperBounds = upperBounds;
+      this.dataSequenceNumber = dataSequenceNumber;
+    }
+
+    @Override
+    public Long pos() {
+      return 0L;
+    }
+
+    @Override
+    public int specId() {
+      return 0;
+    }
+
+    @Override
+    public FileContent content() {
+      return content;
+    }
+
+    @Override
+    public CharSequence path() {
+      return path;
+    }
+
+    @Override
+    public FileFormat format() {
+      return format;
+    }
+
+    @Override
+    public StructLike partition() {
+      return null;
+    }
+
+    @Override
+    public long recordCount() {
+      return recordCount;
+    }
+
+    @Override
+    public long fileSizeInBytes() {
+      return fileSizeInBytes;
+    }
+
+    @Override
+    public Map<Integer, Long> columnSizes() {
+      return null;
+    }
+
+    @Override
+    public Map<Integer, Long> valueCounts() {
+      return null;
+    }
+
+    @Override
+    public Map<Integer, Long> nullValueCounts() {
+      return null;
+    }
+
+    @Override
+    public Map<Integer, Long> nanValueCounts() {
+      return null;
+    }
+
+    @Override
+    public Map<Integer, ByteBuffer> lowerBounds() {
+      return lowerBounds;
+    }
+
+    @Override
+    public Map<Integer, ByteBuffer> upperBounds() {
+      return upperBounds;
+    }
+
+    @Override
+    public ByteBuffer keyMetadata() {
+      return null;
+    }
+
+    @Override
+    public List<Integer> equalityFieldIds() {
+      return equalityFieldIds;
+    }
+
+    @Override
+    public Long dataSequenceNumber() {
+      return dataSequenceNumber;
+    }
+
+    @Override
+    public DeleteFile copy() {
+      return this;
+    }
+
+    @Override
+    public DeleteFile copyWithoutStats() {
+      return this;
+    }
+  }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to