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]