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 c8c0ddca4d4 branch-4.1: [fix](iceberg) keep identity partition values 
of old specs after evolving to unpartitioned(#68124) (#68259)
c8c0ddca4d4 is described below

commit c8c0ddca4d4d452bcb236f2499c5fed9a3c68780
Author: daidai <[email protected]>
AuthorDate: Thu Oct 8 19:10:48 2026 +0800

    branch-4.1: [fix](iceberg) keep identity partition values of old specs 
after evolving to unpartitioned(#68124) (#68259)
    
    Cherry-picked from #68124.
---
 .../doris/datasource/iceberg/IcebergUtils.java     | 14 +++++
 .../datasource/iceberg/source/IcebergScanNode.java | 56 +++++++++++++----
 .../doris/datasource/iceberg/IcebergUtilsTest.java | 19 ++++++
 .../iceberg/source/IcebergScanNodeTest.java        | 73 ++++++++++++++++++++++
 4 files changed, 151 insertions(+), 11 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
index 92fc7744060..dc775bb948c 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
@@ -982,6 +982,20 @@ public class IcebergUtils {
         return partitionInfoMap;
     }
 
+    /**
+     * Whether <b>any</b> partition spec of the table is partitioned. Scan 
planning must use this instead of
+     * the current default spec: after evolving to an unpartitioned spec, data 
files written under an older
+     * partitioned spec still carry their partition metadata (identity values, 
spec id, partition data).
+     */
+    public static boolean hasPartitionedSpec(Table table) {
+        for (PartitionSpec spec : table.specs().values()) {
+            if (spec.isPartitioned()) {
+                return true;
+            }
+        }
+        return false;
+    }
+
     public static List<String> getIdentityPartitionColumns(Table table) {
         return getIdentityPartitionColumns(table, false, false);
     }
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 85f51f07d2c..f8a7f00a23d 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
@@ -199,8 +199,15 @@ public class IcebergScanNode extends FileQueryScanNode {
     private List<String> orderedPartitionMetadataKeys;
     private boolean enableMappingVarbinaryForPartitionMetadata;
     private boolean enableMappingTimestampTzForPartitionMetadata;
+    // Whether the table's CURRENT default spec is partitioned. 
Display/accounting only: it gates what
+    // counts as a scanned partition (selectedPartitionNum, 
numApproximateSplits), preserving the legacy
+    // numbers for a table that evolved to an unpartitioned spec.
     private boolean isPartitionedTable;
     private Boolean requiresRowId;
+    // Whether ANY spec of the table is partitioned. This is the READ gate: 
data files written under an
+    // older partitioned spec still carry their identity values, and must keep 
sending them even after
+    // the default spec evolved to unpartitioned.
+    private boolean hasPartitionedSpec;
     private int formatVersion;
     private ExecutionAuthenticator preExecutionAuthenticator;
     private IcebergRuntimeContext runtimeContext;
@@ -318,6 +325,7 @@ public class IcebergScanNode extends FileQueryScanNode {
             partitionMapInfos = new HashMap<>();
             initializePartitionMetadata();
             isPartitionedTable = icebergTable.spec().isPartitioned();
+            hasPartitionedSpec = IcebergUtils.hasPartitionedSpec(icebergTable);
             // Metadata tables (system tables) are not BaseTable instances, so 
we need to handle this case
             if (icebergTable instanceof BaseTable) {
                 formatVersion = ((BaseTable) 
icebergTable).operations().current().formatVersion();
@@ -510,18 +518,26 @@ public class IcebergScanNode extends FileQueryScanNode {
         }
         enableMappingVarbinaryForPartitionMetadata = 
getEnableMappingVarbinary();
         enableMappingTimestampTzForPartitionMetadata = 
getEnableMappingTimestampTz();
-        orderedPathPartitionKeys = Collections.unmodifiableList(
-                IcebergUtils.getCommonIdentityPartitionColumns(icebergTable,
-                        enableMappingVarbinaryForPartitionMetadata,
-                        enableMappingTimestampTzForPartitionMetadata));
         if (sessionVariable.enableFileScannerV2) {
-            orderedPartitionMetadataKeys = Collections.unmodifiableList(
+            // Classify every identity column of EVERY spec as a partition 
key, like master does. A file
+            // written under an older identity spec carries its value in the 
manifest only, and BE uses a
+            // split partition value just for a column FE marked as a 
partition key, so the intersection
+            // below would leave that value unused and the column read as 
NULL. Files that carry no value
+            // for the column (a spec without it) are unaffected: the v2 
column mapper falls back to the
+            // physical field when the split has no partition value for a 
partition key.
+            orderedPathPartitionKeys = Collections.unmodifiableList(
                     IcebergUtils.getIdentityPartitionColumns(icebergTable,
                             enableMappingVarbinaryForPartitionMetadata,
                             enableMappingTimestampTzForPartitionMetadata));
         } else {
-            orderedPartitionMetadataKeys = orderedPathPartitionKeys;
+            // v1 has no such per-file fallback: a partition key it cannot 
fill from the split is not read
+            // from the file either. Keep the columns common to all specs so 
those files stay readable.
+            orderedPathPartitionKeys = Collections.unmodifiableList(
+                    
IcebergUtils.getCommonIdentityPartitionColumns(icebergTable,
+                            enableMappingVarbinaryForPartitionMetadata,
+                            enableMappingTimestampTzForPartitionMetadata));
         }
+        orderedPartitionMetadataKeys = orderedPathPartitionKeys;
     }
 
     @VisibleForTesting
@@ -2616,7 +2632,10 @@ public class IcebergScanNode extends FileQueryScanNode {
         split.setPartitionSpecId(specId);
         PartitionData partitionData = (PartitionData) dataFile.partition();
         boolean includePartitionData = requiresRowId();
-        if (partitionData != null && (includePartitionData || 
isPartitionedTable)) {
+        // Gate identity values on the table's spec HISTORY, not on the 
current default spec: a file written
+        // under an older identity spec keeps its partition values in the 
manifest, and they are the only
+        // source for an identity partition column that the physical file does 
not store.
+        if (partitionData != null && (includePartitionData || 
hasPartitionedSpec)) {
             PartitionSpec partitionSpec = icebergTable.specs().get(specId);
             Preconditions.checkNotNull(partitionSpec, "Partition spec with 
specId %s not found for table %s",
                     specId, icebergTable.name());
@@ -2631,7 +2650,10 @@ public class IcebergScanNode extends FileQueryScanNode {
                             + specId + ". Rewrite historical data files before 
DELETE or UPDATE.", e);
                 }
             }
-            if (isPartitionedTable) {
+            // The cache doubles as the scanned-partition count only for a 
partitioned current spec. Otherwise
+            // cache only specs with identity fields, or transformed-only old 
partitions would grow it uselessly.
+            if (isPartitionedTable || partitionSpec.fields().stream()
+                    .anyMatch(field -> field.transform().isIdentity())) {
                 Map<String, String> partitionInfoMap = 
partitionMapInfos.computeIfAbsent(
                         Pair.of(specId, partitionData), k -> 
IcebergUtils.getIdentityPartitionInfoMap(
                                 partitionData, partitionSpec, icebergTable, 
sessionVariable.getTimeZone(),
@@ -2813,7 +2835,7 @@ public class IcebergScanNode extends FileQueryScanNode {
             for (FileScanTask task : customFileScanTasks) {
                 splits.add(createIcebergSplit(task));
             }
-            selectedPartitionNum = partitionMapInfos.size();
+            selectedPartitionNum = scannedPartitionNum();
             recordManifestCacheProfile();
             return splits;
         }
@@ -2850,7 +2872,7 @@ public class IcebergScanNode extends FileQueryScanNode {
             }
         }
 
-        selectedPartitionNum = partitionMapInfos.size();
+        selectedPartitionNum = scannedPartitionNum();
         recordManifestCacheProfile();
         return splits;
     }
@@ -3233,9 +3255,21 @@ public class IcebergScanNode extends FileQueryScanNode {
         ((IcebergSplit) splits.get(size - 
1)).setTableLevelRowCount(countPerSplit + totalCount % size);
     }
 
+    /**
+     * The number of scanned partitions reported to EXPLAIN ({@code 
partition=N/M}), the
+     * {@code sql_block_rule} {@code partition_num} guard and the batch-mode 
split estimate. A table whose
+     * CURRENT spec is unpartitioned reports no partitions at all, so the 
partitions that files of an older
+     * spec still carry must not start being counted — {@code 
partitionMapInfos} now collects them because
+     * it doubles as the per-(spec, partition) identity-value cache for the 
read path.
+     */
+    private int scannedPartitionNum() {
+        return isPartitionedTable ? partitionMapInfos.size() : 0;
+    }
+
     @Override
     public int numApproximateSplits() {
-        return NUM_SPLITS_PER_PARTITION * partitionMapInfos.size() > 0 ? 
partitionMapInfos.size() : 1;
+        int partitions = scannedPartitionNum();
+        return NUM_SPLITS_PER_PARTITION * partitions > 0 ? partitions : 1;
     }
 
     private Optional<NotSupportedException> 
checkNotSupportedException(Exception e) {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
index 3717bbc2af3..a9f31393d93 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
@@ -427,6 +427,25 @@ public class IcebergUtilsTest {
         Assert.assertTrue(exception.getMessage().contains("db.tbl"));
     }
 
+    @Test
+    public void testHasPartitionedSpecLooksAtEverySpecNotOnlyTheCurrentOne() {
+        Schema schema = new Schema(
+                Types.NestedField.required(1, "id", Types.IntegerType.get()),
+                Types.NestedField.required(2, "p", Types.IntegerType.get()));
+        PartitionSpec identity = 
PartitionSpec.builderFor(schema).identity("p").withSpecId(0).build();
+        PartitionSpec unpartitioned = PartitionSpec.unpartitioned();
+
+        // A table that evolved identity(p) -> unpartitioned still has files 
carrying p in their
+        // partition metadata, so scan planning must keep sending those values.
+        Table evolved = Mockito.mock(Table.class);
+        Mockito.when(evolved.specs()).thenReturn(ImmutableMap.of(0, identity, 
1, unpartitioned));
+        Assert.assertTrue(IcebergUtils.hasPartitionedSpec(evolved));
+
+        Table neverPartitioned = Mockito.mock(Table.class);
+        Mockito.when(neverPartitioned.specs()).thenReturn(ImmutableMap.of(0, 
unpartitioned));
+        Assert.assertFalse(IcebergUtils.hasPartitionedSpec(neverPartitioned));
+    }
+
     @Test
     public void testEmptyNameMappingStillParsesAsAuthoritativeMapping() throws 
Exception {
         Table table = Mockito.mock(Table.class);
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 3d02dd83341..84d5f625d53 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
@@ -1824,6 +1824,69 @@ public class IcebergScanNodeTest {
                 IcebergUtils.getSerializedInitialDefaults(fields, 
false).get(7));
     }
 
+    @Test
+    public void 
testSplitKeepsOldSpecIdentityValuesAfterEvolvingToUnpartitioned() throws 
Exception {
+        // A file written under identity(p) keeps p=7 only in its manifest 
partition metadata. Gating on
+        // the CURRENT spec drops that value once the default spec evolves to 
unpartitioned, and keying
+        // the partition-key classification off the columns common to all 
specs leaves BE reading p from
+        // a file that does not store it, so the column comes back NULL in 
both cases.
+        Schema schema = new Schema(
+                Types.NestedField.required(1, "id", Types.IntegerType.get()),
+                Types.NestedField.required(2, "p", Types.IntegerType.get()));
+        PartitionSpec spec = 
PartitionSpec.builderFor(schema).identity("p").build();
+        String location = 
temporaryFolder.newFolder("identity_evolved_table").toURI().toString();
+        Table table = new HadoopTables(new Configuration()).create(schema, 
spec, location);
+        table.newFastAppend()
+                .appendFile(DataFiles.builder(table.spec())
+                        .withPath(location + "data/p=7/old.parquet")
+                        .withFileSizeInBytes(1024L)
+                        .withRecordCount(2L)
+                        .withFormat(FileFormat.PARQUET)
+                        .withPartitionPath("p=7")
+                        .build())
+                .commit();
+        table.updateSpec().removeField("p").commit();
+        Assert.assertTrue(table.spec().isUnpartitioned());
+
+        SessionVariable sessionVariable = new SessionVariable();
+        sessionVariable.enableFileScannerV2 = true;
+        TestIcebergScanNode node = new TestIcebergScanNode(sessionVariable);
+        setIcebergTable(node, table);
+        setPrivateField(node, "formatVersion", 2);
+        setPrivateField(node, "storagePropertiesMap", Collections.emptyMap());
+        setPrivateField(node, "partitionMapInfos", new HashMap<>());
+        // Exactly what doInitialize computes for this table.
+        setPrivateField(node, "isPartitionedTable", 
table.spec().isPartitioned());
+        setPrivateField(node, "hasPartitionedSpec", 
IcebergUtils.hasPartitionedSpec(table));
+
+        // p must stay a partition key even though the current spec no longer 
mentions it; otherwise BE
+        // ignores the split value and reads the (absent) physical column.
+        Assert.assertEquals(ImmutableList.of("p"), 
orderedPathPartitionKeys(node));
+
+        FileScanTask task;
+        try (CloseableIterable<FileScanTask> tasks = 
table.newScan().planFiles()) {
+            Iterator<FileScanTask> iterator = tasks.iterator();
+            Assert.assertTrue(iterator.hasNext());
+            task = iterator.next();
+        }
+
+        IcebergSplit split = createIcebergSplit(node, task);
+        Assert.assertEquals(ImmutableMap.of("p", "7"), 
split.getIcebergPartitionValues());
+        Assert.assertEquals(Integer.valueOf(0), split.getPartitionSpecId());
+        // A plain read does not project the row ID, so no row-ID partition 
JSON is generated.
+        Assert.assertNull(split.getPartitionDataJson());
+
+        // What BE actually receives: p=7 as a columns-from-path constant for 
this range.
+        TFileRangeDesc rangeDesc = new TFileRangeDesc();
+        setIcebergParams(node, rangeDesc, split);
+        Assert.assertEquals(ImmutableList.of("p"), 
rangeDesc.getColumnsFromPathKeys());
+        Assert.assertEquals(ImmutableList.of("7"), 
rangeDesc.getColumnsFromPath());
+
+        // Display parity: the CURRENT spec is unpartitioned, so the table 
still reports no scanned
+        // partitions and the batch-mode estimate stays at the unpartitioned 
default.
+        Assert.assertEquals(1, node.numApproximateSplits());
+    }
+
     @Test
     public void testBatchSplitCarriesDroppedEqualitySchemaThroughThrift() 
throws Exception {
         Types.NestedField id = Types.NestedField.required(1, "id", 
Types.LongType.get());
@@ -1958,6 +2021,7 @@ public class IcebergScanNodeTest {
             node.addSlot(0, new Column("record_key", Type.INT));
             setIcebergTable(node, table);
             setPrivateField(node, "isPartitionedTable", keepPartitioned);
+            setPrivateField(node, "hasPartitionedSpec", 
IcebergUtils.hasPartitionedSpec(table));
             setPrivateField(node, "storagePropertiesMap", 
Collections.emptyMap());
             setPrivateField(node, "formatVersion", 2);
             setPrivateField(node, "partitionMapInfos", new HashMap<>());
@@ -2043,6 +2107,7 @@ public class IcebergScanNodeTest {
         TestIcebergScanNode node = new TestIcebergScanNode(new 
SessionVariable());
         setIcebergTable(node, table);
         setPrivateField(node, "isPartitionedTable", keepPartitioned);
+        setPrivateField(node, "hasPartitionedSpec", 
IcebergUtils.hasPartitionedSpec(table));
         setPrivateField(node, "storagePropertiesMap", Collections.emptyMap());
         setPrivateField(node, "formatVersion", 2);
         setPrivateField(node, "partitionMapInfos", new HashMap<>());
@@ -2181,6 +2246,7 @@ public class IcebergScanNodeTest {
         node.addSlot(0, IcebergRowId.createHiddenColumn());
         setIcebergTable(node, table);
         setPrivateField(node, "isPartitionedTable", 
table.spec().isPartitioned());
+        setPrivateField(node, "hasPartitionedSpec", 
IcebergUtils.hasPartitionedSpec(table));
         setPrivateField(node, "storagePropertiesMap", Collections.emptyMap());
         setPrivateField(node, "formatVersion", 2);
         setPrivateField(node, "partitionMapInfos", new HashMap<>());
@@ -4097,6 +4163,13 @@ public class IcebergScanNodeTest {
         return task;
     }
 
+    @SuppressWarnings("unchecked")
+    private static List<String> orderedPathPartitionKeys(IcebergScanNode node) 
throws Exception {
+        Method method = 
IcebergScanNode.class.getDeclaredMethod("getOrderedPathPartitionKeys");
+        method.setAccessible(true);
+        return (List<String>) method.invoke(node);
+    }
+
     private static IcebergSplit createIcebergSplit(IcebergScanNode node, 
FileScanTask task)
             throws Exception {
         Method method = 
IcebergScanNode.class.getDeclaredMethod("createIcebergSplit", 
FileScanTask.class);


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

Reply via email to