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]