This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 2f84f2fb535 branch-4.1: [fix](iceberg) Resolve dropped equality-delete
fields in manifest-cache planning (#68433)
2f84f2fb535 is described below
commit 2f84f2fb5354eb6af2fb572779834b632e614f59
Author: daidai <[email protected]>
AuthorDate: Thu Oct 8 18:13:53 2026 +0800
branch-4.1: [fix](iceberg) Resolve dropped equality-delete fields in
manifest-cache planning (#68433)
### What problem does this PR solve?
Issue Number: None
Related PR: #67479, #67687
Problem Summary:
#67479 lets `DeleteFileIndex` resolve equality-delete fields that were
later dropped, by looking them up in historical schemas. It covers the
default planner but misses the manifest-cache planner in
`IcebergScanNode`, which still passes only `specsById`. When the
manifest cache is enabled, planning fails with `Cannot find field for ID
N` and falls back to the SDK scan: results are correct, but each query
logs a WARN and skips the cache.
The master port (#67687) already covers this path. This PR applies the
same one-line change to branch-4.1: pass `icebergTable.schemas()` to the
index.
### Release note
None
### Check List (For Author)
- Test
- [ ] Regression test
- [x] Unit Test
- [x] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason
Manual test: with the manifest cache enabled, `EXPLAIN VERBOSE` shows
`failures=0` instead of `1`, and the fallback WARN is gone.
- Behavior changed:
- [x] No.
- [ ] Yes.
- Does this need documentation?
- [x] No.
- [ ] Yes.
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label
---
.../datasource/iceberg/source/IcebergScanNode.java | 11 ++-
.../java/org/apache/iceberg/DeleteFileIndex.java | 3 +-
.../iceberg/source/IcebergScanNodeTest.java | 84 ++++++++++++++++++++++
3 files changed, 94 insertions(+), 4 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
index de8ce891662..85f51f07d2c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
@@ -2440,10 +2440,15 @@ public class IcebergScanNode extends FileQueryScanNode {
}
// Build delete file index for efficient lookup of deletes applicable
to each data file
- DeleteFileIndex deleteIndex = DeleteFileIndex.builderFor(deleteFiles)
+ DeleteFileIndex.Builder deleteIndexBuilder =
DeleteFileIndex.builderFor(deleteFiles)
.specsById(specsById)
- .caseSensitive(caseSensitive)
- .build();
+ .caseSensitive(caseSensitive);
+ // Equality deletes may reference fields dropped from the current
schema, so resolve their
+ // field IDs against every historical schema. Other scans skip
indexing the schema history.
+ if (deleteFiles.stream().anyMatch(file -> file.content() ==
FileContent.EQUALITY_DELETES)) {
+ deleteIndexBuilder.schemasById(icebergTable.schemas());
+ }
+ DeleteFileIndex deleteIndex = deleteIndexBuilder.build();
// ========== Phase 2: Load data files and create scan tasks ==========
List<FileScanTask> tasks = new ArrayList<>();
diff --git a/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
b/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
index 630f19c81a4..da7f7803835 100644
--- a/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
+++ b/fe/fe-core/src/main/java/org/apache/iceberg/DeleteFileIndex.java
@@ -400,7 +400,8 @@ public class DeleteFileIndex {
return this;
}
- Builder schemasById(Map<Integer, Schema> newSchemasById) {
+ // changed to public method: equality deletes may reference fields dropped
from the current schema.
+ public Builder schemasById(Map<Integer, Schema> newSchemasById) {
this.schemasById = newSchemasById;
return this;
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
index 3e17f77ac74..3d02dd83341 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
@@ -30,6 +30,7 @@ import org.apache.doris.analysis.TupleDescriptor;
import org.apache.doris.analysis.TupleId;
import org.apache.doris.catalog.Column;
import org.apache.doris.catalog.DatabaseIf;
+import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.StructField;
import org.apache.doris.catalog.StructType;
import org.apache.doris.catalog.TableIf;
@@ -40,7 +41,10 @@ import
org.apache.doris.common.security.authentication.ExecutionAuthenticator;
import
org.apache.doris.common.security.authentication.HadoopExecutionAuthenticator;
import org.apache.doris.common.util.LocationPath;
import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.ExternalCatalog;
+import org.apache.doris.datasource.ExternalMetaCacheMgr;
import org.apache.doris.datasource.ExternalScanNode;
+import org.apache.doris.datasource.ExternalTable;
import org.apache.doris.datasource.FederationBackendPolicy;
import org.apache.doris.datasource.SplitAssignment;
import org.apache.doris.datasource.SplitGenerator;
@@ -57,6 +61,8 @@ import
org.apache.doris.datasource.iceberg.IcebergSnapshotCacheValue;
import org.apache.doris.datasource.iceberg.IcebergSysExternalTable;
import org.apache.doris.datasource.iceberg.IcebergTableCacheValue;
import org.apache.doris.datasource.iceberg.IcebergUtils;
+import org.apache.doris.datasource.iceberg.cache.IcebergManifestCacheLoader;
+import org.apache.doris.datasource.iceberg.cache.ManifestCacheValue;
import org.apache.doris.datasource.iceberg.helper.IcebergWriterHelper;
import org.apache.doris.datasource.mvcc.MvccTableInfo;
import org.apache.doris.nereids.StatementContext;
@@ -97,6 +103,7 @@ import org.apache.iceberg.FileContent;
import org.apache.iceberg.FileFormat;
import org.apache.iceberg.FileMetadata;
import org.apache.iceberg.FileScanTask;
+import org.apache.iceberg.ManifestFiles;
import org.apache.iceberg.MetadataColumns;
import org.apache.iceberg.MetadataTableType;
import org.apache.iceberg.PartitionData;
@@ -125,6 +132,7 @@ import org.junit.Assert;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.TemporaryFolder;
+import org.mockito.MockedStatic;
import org.mockito.Mockito;
import java.io.Closeable;
@@ -3594,6 +3602,82 @@ public class IcebergScanNodeTest {
}
}
+ @Test
+ public void testManifestCachePlanningResolvesDroppedEqualityDeleteField()
throws Exception {
+ Schema schema = new Schema(
+ Types.NestedField.optional(1, "k", Types.IntegerType.get()),
+ Types.NestedField.optional(2, "row_tag",
Types.StringType.get()));
+ String tableLocation = temporaryFolder.getRoot().toPath()
+
.resolve("manifest_cache_dropped_equality_key").toUri().toString();
+ Table table = new HadoopTables(new Configuration()).create(
+ schema, PartitionSpec.unpartitioned(), SortOrder.unsorted(),
+ ImmutableMap.of(TableProperties.FORMAT_VERSION, "2"),
tableLocation);
+ table.newFastAppend().appendFile(DataFiles.builder(table.spec())
+ .withPath(tableLocation + "/data/old.parquet")
+ .withFormat(FileFormat.PARQUET)
+ .withFileSizeInBytes(10)
+ .withRecordCount(2)
+ .build()).commit();
+
table.newRowDelta().addDeletes(FileMetadata.deleteFileBuilder(table.spec())
+ .ofEqualityDeletes(1)
+ .withPath(tableLocation + "/data/eq-delete.parquet")
+ .withFormat(FileFormat.PARQUET)
+ .withFileSizeInBytes(10)
+ .withRecordCount(1)
+ .build()).commit();
+ // The live equality delete still references field 1 after it leaves
the current schema.
+ table.updateSchema().deleteColumn("k").commit();
+ table.newFastAppend().appendFile(DataFiles.builder(table.spec())
+ .withPath(tableLocation + "/data/new.parquet")
+ .withFormat(FileFormat.PARQUET)
+ .withFileSizeInBytes(10)
+ .withRecordCount(1)
+ .build()).commit();
+
+ TestIcebergScanNode node = new TestIcebergScanNode(new
SessionVariable());
+ IcebergSource source = Mockito.mock(IcebergSource.class);
+
Mockito.when(source.getCatalog()).thenReturn(Mockito.mock(ExternalCatalog.class));
+
Mockito.when(source.getTargetTable()).thenReturn(Mockito.mock(ExternalTable.class));
+ setPrivateField(node, "source", source);
+ setPrivateField(node, "icebergTable", table);
+ Env env = Mockito.mock(Env.class);
+
Mockito.when(env.getExtMetaCacheMgr()).thenReturn(Mockito.mock(ExternalMetaCacheMgr.class));
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class);
+ MockedStatic<IcebergManifestCacheLoader> loader =
+ Mockito.mockStatic(IcebergManifestCacheLoader.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+ loader.when(() ->
IcebergManifestCacheLoader.loadDeleteFilesWithCache(Mockito.any(),
Mockito.any(),
+ Mockito.any(), Mockito.any(), Mockito.any(),
Mockito.any()))
+ .thenAnswer(invocation -> {
+ List<DeleteFile> files = new ArrayList<>();
+ try (CloseableIterable<DeleteFile> reader =
ManifestFiles.readDeleteManifest(
+ invocation.getArgument(2), table.io(),
table.specs())) {
+ reader.forEach(file -> files.add(file.copy()));
+ }
+ return ManifestCacheValue.forDeleteFiles(files);
+ });
+ loader.when(() ->
IcebergManifestCacheLoader.loadDataFilesWithCache(Mockito.any(), Mockito.any(),
+ Mockito.any(), Mockito.any(), Mockito.any(),
Mockito.any()))
+ .thenAnswer(invocation -> {
+ List<DataFile> files = new ArrayList<>();
+ try (CloseableIterable<DataFile> reader =
ManifestFiles.read(
+ invocation.getArgument(2), table.io(),
table.specs())) {
+ reader.forEach(file -> files.add(file.copy()));
+ }
+ return ManifestCacheValue.forDataFiles(files);
+ });
+
+ List<FileScanTask> tasks = node.loadFileScanTasksWithManifestCache(
+ table.newScan(), table.currentSnapshot());
+
+ Map<String, Integer> deleteCounts = new HashMap<>();
+ tasks.forEach(task -> deleteCounts.put(task.file().location(),
task.deletes().size()));
+ Assert.assertEquals(ImmutableMap.of(
+ tableLocation + "/data/old.parquet", 1,
+ tableLocation + "/data/new.parquet", 0), deleteCounts);
+ }
+ }
+
@Test
public void testPinnedBranchUsesFrozenSnapshotWithCurrentSchema() throws
Exception {
Schema snapshotSchema = new Schema(11, ImmutableList.of(
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]