This is an automated email from the ASF dual-hosted git repository. yiguolei pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/doris.git
commit 97df025c483310a9a6674f7c3b2587b730fa5557 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]
