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 502268d05a3 [fix](iceberg) Preserve file specs after evolving to 
unpartitioned tables (#68481)
502268d05a3 is described below

commit 502268d05a3a12f7b296a38b2c6986f0a59f826e
Author: Gabriel <[email protected]>
AuthorDate: Tue Sep 29 15:37:59 2026 +0800

    [fix](iceberg) Preserve file specs after evolving to unpartitioned tables 
(#68481)
    
    ### What problem does this PR solve?
    
    On `branch-4.1`, Iceberg scan splits only carry partition metadata when
    the table's current spec is partitioned. After an external REPLACE
    changes a partitioned table to an unpartitioned table, the old spec
    remains in the table metadata but new files use a different spec ID.
    Missing split metadata makes the BE default to spec 0, so DELETE can
    fail with `No partition data for partitioned table` when committing the
    position-delete files. Dropping the last partition field also affects
    tables containing files from both specs.
    
    Preserve each data file's spec ID regardless of the current table spec.
    Generate partition JSON only for scans projecting the hidden row ID,
    with the projection requirement cached once per scan. Partition-pruning
    metadata remains independent of this write payload.
    
    Historical BINARY/FIXED partition values round-trip through the target
    branch's lossless `0x`-prefixed hexadecimal representation, with
    matching decoding and fixed-length validation. UUID/TIME values also
    round-trip through delete metadata. Unsupported partition values are
    never serialized as NULL. Ordinary reads omit unavailable row-ID
    partition metadata, while scans that materialize the hidden row ID
    reject unsupported partition types before generating delete metadata.
    This keeps historical files readable after their partition source column
    is dropped. Non-batch split planning preserves the actionable
    UserException even when Hadoop authentication wraps it, including the
    spec ID and rewrite guidance. Negative fractional timestamp values use
    floor-based seconds/nanoseconds normalization, and TIMESTAMPTZ transport
    retains an explicit UTC offset across DST overlaps.
    
    ### Release note
    
    Fix Iceberg row-level deletes after replacing a partitioned table with
    an unpartitioned table or dropping its last partition field. Keep
    historical files readable after dropping a partition source column, and
    reject row-ID-dependent operations when the old partition type is
    unavailable.
    
    ### Tests
    
    - Added two unit tests using real Iceberg metadata transitions. Both
    fail on the original code because the spec ID is missing from the
    serialized scan range, and pass with the fix. They also validate
    conversion to position-delete metadata for the file's own spec.
    - Added dropped-source-column tests using real Iceberg metadata,
    covering NULL/non-NULL historical values and both currently partitioned
    and unpartitioned tables. They verify readable scan ranges without
    fabricated partition JSON and explicit rejection when row IDs are
    required.
    - Added ordinary-read checks that partition JSON is absent while each
    file retains its spec ID and applicable partition-pruning values.
    Dropped-source-column error tests now call the public getSplits entry
    point, including a real Hadoop doAs wrapper, and verify the
    UserException type, message, spec ID, guidance, and original cause.
    Existing delete metadata tests explicitly project the hidden row ID.
    - Added six evolved-spec tests for BINARY, FIXED, UUID, TIME, TIMESTAMP,
    and TIMESTAMPTZ. These reproduced the review findings before the
    follow-up fix. Coverage includes empty/NULL binary values, fixed-length
    validation, read-only/direct buffer bounds, and negative fractional
    timestamps.
    - Ran `IcebergScanNodeTest`, `IcebergUtilsTest`,
    `IcebergWriterHelperTest`, and `IcebergTransactionTest`: **217 tests
    passed**, no failures or skips.
    - FE Checkstyle: passed during the test build with
    `-Dcheckstyle.skip=false`.
    - Added `test_iceberg_delete_unpartitioned_evolution` to the external
    Iceberg regression suite, covering Spark REPLACE followed by repeated
    Doris DELETE, plus mixed-spec DELETE/UPDATE after dropping the last
    partition field. It checks both query results against Spark and
    delete-file spec IDs, with Parquet and ORC coverage.
    - Extended the external suite with historical binary-partition DELETE,
    pre-epoch timestamp-partition DELETE/UPDATE, and Parquet/ORC
    dropped-source-column reads with safe DELETE/UPDATE rejection and
    unchanged query results.
    - Regression Groovy syntax check passed. The external regression suite
    has **not been run end-to-end locally**; it requires the Iceberg/Spark
    test environment. Spark REPLACE exercises the metadata transition
    reported with Trino RTAS.
    
    ### Check List (For Author)
    
    - Test
        - [x] Regression test
        - [x] Unit Test
    - Behavior changed:
    - [x] Yes. Preserve per-file partition metadata for evolved
    unpartitioned tables.
    - Does this need documentation?
        - [x] No.
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Confirm document
    - [ ] Add branch pick label
---
 .../doris/datasource/iceberg/IcebergUtils.java     |  11 +-
 .../datasource/iceberg/source/IcebergScanNode.java |  46 ++-
 .../doris/datasource/iceberg/IcebergUtilsTest.java |  23 ++
 .../iceberg/source/IcebergScanNodeTest.java        | 339 +++++++++++++++++++++
 ...t_iceberg_delete_unpartitioned_evolution.groovy | 198 ++++++++++++
 5 files changed, 601 insertions(+), 16 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 bd5d2c8cb65..92fc7744060 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
@@ -1104,13 +1104,8 @@ public class IcebergUtils {
         for (int i = 0; i < fields.size(); i++) {
             NestedField field = fields.get(i);
             Object value = partitionData.get(i);
-            try {
-                partitionValues.add(serializePartitionValue(field.type(), 
value, timeZone));
-            } catch (UnsupportedOperationException e) {
-                LOG.warn("Failed to serialize Iceberg partition value for 
field {}: {}", field.name(),
-                        e.getMessage());
-                partitionValues.add(null);
-            }
+            // These values also identify delete-file partitions; an 
unsupported value must never become NULL.
+            partitionValues.add(serializePartitionValue(field.type(), value, 
timeZone));
         }
         return partitionValues;
     }
@@ -1407,6 +1402,8 @@ public class IcebergUtils {
                     return (int) LocalDate.parse(valueStr, 
DateTimeFormatter.ISO_LOCAL_DATE).toEpochDay();
                 case TIMESTAMP:
                     return parseTimestampToMicros(valueStr, (TimestampType) 
icebergType);
+                case TIME:
+                    return LocalTime.parse(valueStr, 
DateTimeFormatter.ISO_LOCAL_TIME).toNanoOfDay() / 1000;
                 case DECIMAL:
                     return new BigDecimal(valueStr);
                 default:
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 75489d13b37..de8ce891662 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
@@ -200,6 +200,7 @@ public class IcebergScanNode extends FileQueryScanNode {
     private boolean enableMappingVarbinaryForPartitionMetadata;
     private boolean enableMappingTimestampTzForPartitionMetadata;
     private boolean isPartitionedTable;
+    private Boolean requiresRowId;
     private int formatVersion;
     private ExecutionAuthenticator preExecutionAuthenticator;
     private IcebergRuntimeContext runtimeContext;
@@ -1658,6 +1659,12 @@ public class IcebergScanNode extends FileQueryScanNode {
         try {
             return preExecutionAuthenticator.execute(() -> 
doGetSplits(numBackends));
         } catch (Exception e) {
+            // Authentication can wrap user errors; keep their guidance 
instead of only the root cause.
+            for (Throwable cause : ExceptionUtils.getThrowableList(e)) {
+                if (cause instanceof UserException) {
+                    throw (UserException) cause;
+                }
+            }
             Optional<NotSupportedException> opt = 
checkNotSupportedException(e);
             if (opt.isPresent()) {
                 throw opt.get();
@@ -2556,6 +2563,15 @@ public class IcebergScanNode extends FileQueryScanNode {
         return LocationPath.of(path, storagePropertiesMap);
     }
 
+    private boolean requiresRowId() {
+        if (requiresRowId == null) {
+            // Slots are finalized by split planning; inspect them once for 
all files in this scan.
+            requiresRowId = desc.getSlots().stream().anyMatch(slot -> 
slot.getColumn() != null
+                    && 
Column.ICEBERG_ROWID_COL.equalsIgnoreCase(slot.getColumn().getName()));
+        }
+        return requiresRowId;
+    }
+
     private Split createIcebergSplit(FileScanTask fileScanTask) throws 
UserException {
         DataFile dataFile = fileScanTask.file();
         String originalPath = dataFile.path().toString();
@@ -2589,16 +2605,28 @@ public class IcebergScanNode extends FileQueryScanNode {
         }
         split.setTableFormatType(TableFormatType.ICEBERG);
         split.setTargetSplitSize(selectFeSplitSize(fileScanTask, 
targetSplitSize));
-        if (isPartitionedTable) {
-            int specId = fileScanTask.file().specId();
+        // REPLACE or partition evolution can leave an unpartitioned table 
with historical specs.
+        // Row-level deletes must retain each file's spec instead of 
defaulting to historical spec 0.
+        int specId = dataFile.specId();
+        split.setPartitionSpecId(specId);
+        PartitionData partitionData = (PartitionData) dataFile.partition();
+        boolean includePartitionData = requiresRowId();
+        if (partitionData != null && (includePartitionData || 
isPartitionedTable)) {
             PartitionSpec partitionSpec = icebergTable.specs().get(specId);
             Preconditions.checkNotNull(partitionSpec, "Partition spec with 
specId %s not found for table %s",
                     specId, icebergTable.name());
-            PartitionData partitionData = (PartitionData) 
fileScanTask.file().partition();
-            if (partitionData != null) {
-                split.setPartitionSpecId(specId);
-                split.setPartitionDataJson(IcebergUtils.getPartitionDataJson(
-                        partitionData, partitionSpec, 
sessionVariable.getTimeZone()));
+            // Only row-ID consumers need this JSON; ordinary reads still 
retain spec and pruning metadata.
+            if (includePartitionData) {
+                try {
+                    
split.setPartitionDataJson(IcebergUtils.getPartitionDataJson(
+                            partitionData, partitionSpec, 
sessionVariable.getTimeZone()));
+                } catch (UnsupportedOperationException e) {
+                    // A dropped source column can leave UNKNOWN; DML must not 
substitute a NULL partition.
+                    throw new UserException("Cannot produce Iceberg row IDs 
with unsupported partition types in spec "
+                            + specId + ". Rewrite historical data files before 
DELETE or UPDATE.", e);
+                }
+            }
+            if (isPartitionedTable) {
                 Map<String, String> partitionInfoMap = 
partitionMapInfos.computeIfAbsent(
                         Pair.of(specId, partitionData), k -> 
IcebergUtils.getIdentityPartitionInfoMap(
                                 partitionData, partitionSpec, icebergTable, 
sessionVariable.getTimeZone(),
@@ -2610,9 +2638,9 @@ public class IcebergScanNode extends FileQueryScanNode {
                 if (!partitionInfoMap.isEmpty()) {
                     split.setIcebergPartitionValues(partitionInfoMap);
                 }
-            } else {
-                partitionMapInfos.put(Pair.of(specId, null), 
Collections.emptyMap());
             }
+        } else if (partitionData == null && isPartitionedTable) {
+            partitionMapInfos.put(Pair.of(specId, null), 
Collections.emptyMap());
         }
         return split;
     }
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 66bb9fb2c66..3717bbc2af3 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
@@ -1016,6 +1016,29 @@ public class IcebergUtilsTest {
         Assert.assertNull(partitionInfoMap);
     }
 
+    @Test
+    public void testBinaryPartitionJsonPreservesBufferBounds() {
+        for (org.apache.iceberg.types.Type type : 
Arrays.asList(Types.BinaryType.get(), Types.FixedType.ofLength(4))) {
+            Schema schema = new Schema(Types.NestedField.optional(1, 
"partition_key", type));
+            PartitionSpec spec = 
PartitionSpec.builderFor(schema).identity("partition_key").build();
+            ByteBuffer buffer = ByteBuffer.allocateDirect(6);
+            buffer.put(new byte[] {42, 0, (byte) 0xff, (byte) 0x80, 0x2f, 42});
+            buffer.position(1);
+            buffer.limit(5);
+            ByteBuffer value = buffer.asReadOnlyBuffer();
+            PartitionData partition = new PartitionData(spec.partitionType());
+            partition.set(0, value);
+            List<String> encoded = IcebergUtils.parsePartitionValuesFromJson(
+                    IcebergUtils.getPartitionDataJson(partition, spec, "UTC"));
+            Assert.assertEquals(Collections.singletonList("0x00ff802f"), 
encoded);
+            Assert.assertEquals(1, value.position());
+            Assert.assertEquals(5, value.limit());
+            Assert.assertEquals(value, 
IcebergUtils.parsePartitionValueFromString(encoded.get(0), type));
+        }
+        Assert.assertThrows(IllegalArgumentException.class,
+                () -> IcebergUtils.parsePartitionValueFromString("0x00ff802f", 
Types.FixedType.ofLength(3)));
+    }
+
     @Test
     public void testGetIdentityPartitionColumnsIgnoresTransformPartitions() {
         Schema schema = new Schema(
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 a104f5982ce..3e17f77ac74 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
@@ -37,6 +37,7 @@ import org.apache.doris.catalog.Type;
 import org.apache.doris.common.ThreadPoolManager;
 import org.apache.doris.common.UserException;
 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.ExternalScanNode;
@@ -49,12 +50,14 @@ import 
org.apache.doris.datasource.iceberg.IcebergExternalCatalog;
 import org.apache.doris.datasource.iceberg.IcebergExternalTable;
 import org.apache.doris.datasource.iceberg.IcebergMvccSnapshot;
 import org.apache.doris.datasource.iceberg.IcebergPartitionInfo;
+import org.apache.doris.datasource.iceberg.IcebergRowId;
 import org.apache.doris.datasource.iceberg.IcebergRuntimeContext;
 import org.apache.doris.datasource.iceberg.IcebergSnapshot;
 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.helper.IcebergWriterHelper;
 import org.apache.doris.datasource.mvcc.MvccTableInfo;
 import org.apache.doris.nereids.StatementContext;
 import org.apache.doris.planner.PlanNodeId;
@@ -65,10 +68,13 @@ import org.apache.doris.system.Backend;
 import org.apache.doris.thrift.TAccessPathType;
 import org.apache.doris.thrift.TColumnAccessPath;
 import org.apache.doris.thrift.TDataAccessPath;
+import org.apache.doris.thrift.TFileContent;
 import org.apache.doris.thrift.TFileFormatType;
 import org.apache.doris.thrift.TFileRangeDesc;
 import org.apache.doris.thrift.TFileScanRangeParams;
+import org.apache.doris.thrift.TIcebergCommitData;
 import org.apache.doris.thrift.TIcebergDeleteFileDesc;
+import org.apache.doris.thrift.TIcebergFileDesc;
 import org.apache.doris.thrift.TMetaAccessPath;
 import org.apache.doris.thrift.TPushAggOp;
 import org.apache.doris.thrift.schema.external.TField;
@@ -78,6 +84,7 @@ import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.ImmutableSet;
 import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.security.UserGroupInformation;
 import org.apache.iceberg.AppendFiles;
 import org.apache.iceberg.BaseMetadataTable;
 import org.apache.iceberg.BaseTable;
@@ -1880,6 +1887,338 @@ public class IcebergScanNodeTest {
                 .isSetEqualityDeleteSchema());
     }
 
+    @Test
+    public void testDeleteSplitKeepsSpecAfterReplacingPartitionedTable() 
throws Exception {
+        HadoopTables tables = new HadoopTables(new Configuration());
+        Schema schema = new Schema(Types.NestedField.required(1, "record_key", 
Types.IntegerType.get()));
+        String location = 
temporaryFolder.newFolder("replace_partitioned").toURI().toString();
+        Table table = tables.create(schema, 
PartitionSpec.builderFor(schema).identity("record_key").build(),
+                Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"), 
location);
+        
table.newAppend().appendFile(DataFiles.builder(table.spec()).withPath("file:///warehouse/old.parquet")
+                
.withPartitionPath("record_key=7").withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+        org.apache.iceberg.Transaction replace = tables.buildTable(location, 
schema)
+                
.withPartitionSpec(PartitionSpec.unpartitioned()).replaceTransaction();
+        
replace.newAppend().appendFile(DataFiles.builder(replace.table().spec()).withPath("file:///warehouse/new.parquet")
+                .withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+        replace.commitTransaction();
+        table.refresh();
+
+        Assert.assertTrue(table.spec().isUnpartitioned());
+        Assert.assertTrue(table.specs().get(0).isPartitioned());
+        Assert.assertNotEquals(0, table.spec().specId());
+        assertDeleteSplitPartitionMetadata(table);
+    }
+
+    @Test
+    public void 
testDeleteSplitKeepsOldAndNewSpecsAfterDroppingPartitionField() throws 
Exception {
+        Schema schema = new Schema(Types.NestedField.required(1, "record_key", 
Types.IntegerType.get()));
+        Table table = new HadoopTables(new Configuration()).create(schema,
+                
PartitionSpec.builderFor(schema).identity("record_key").build(),
+                Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"),
+                
temporaryFolder.newFolder("drop_partition_field").toURI().toString());
+        
table.newAppend().appendFile(DataFiles.builder(table.spec()).withPath("file:///warehouse/old.parquet")
+                
.withPartitionPath("record_key=7").withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+        table.updateSpec().removeField("record_key").commit();
+        
table.newAppend().appendFile(DataFiles.builder(table.spec()).withPath("file:///warehouse/new.parquet")
+                .withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+
+        Assert.assertTrue(table.spec().isUnpartitioned());
+        assertDeleteSplitPartitionMetadata(table);
+    }
+
+    @Test
+    public void testReadSplitsSkipPartitionJson() throws Exception {
+        for (boolean keepPartitioned : Arrays.asList(false, true)) {
+            Schema schema = new Schema(Types.NestedField.required(1, 
"record_key", Types.IntegerType.get()));
+            Table table = new HadoopTables(new Configuration()).create(schema,
+                    
PartitionSpec.builderFor(schema).identity("record_key").build(),
+                    Collections.singletonMap(TableProperties.FORMAT_VERSION, 
"2"),
+                    temporaryFolder.newFolder().toURI().toString());
+            
table.newAppend().appendFile(DataFiles.builder(table.spec()).withPath("file:///warehouse/old.parquet")
+                    
.withPartitionPath("record_key=7").withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+            if (!keepPartitioned) {
+                table.updateSpec().removeField("record_key").commit();
+            }
+            DataFiles.Builder newFile = 
DataFiles.builder(table.spec()).withPath("file:///warehouse/new.parquet")
+                    .withRecordCount(1).withFileSizeInBytes(128);
+            if (keepPartitioned) {
+                newFile.withPartitionPath("record_key=9");
+            }
+            table.newAppend().appendFile(newFile.build()).commit();
+
+            TestIcebergScanNode node = new TestIcebergScanNode(new 
SessionVariable());
+            node.addSlot(0, new Column("record_key", Type.INT));
+            setIcebergTable(node, table);
+            setPrivateField(node, "isPartitionedTable", keepPartitioned);
+            setPrivateField(node, "storagePropertiesMap", 
Collections.emptyMap());
+            setPrivateField(node, "formatVersion", 2);
+            setPrivateField(node, "partitionMapInfos", new HashMap<>());
+            setPrivateField(node, "orderedPathPartitionKeys", 
Collections.emptyList());
+            setPrivateField(node, "orderedPartitionMetadataKeys", 
Collections.emptyList());
+            try (CloseableIterable<FileScanTask> tasks = 
table.newScan().planFiles()) {
+                int fileCount = 0;
+                for (FileScanTask task : tasks) {
+                    IcebergSplit split = createIcebergSplit(node, task);
+                    TFileRangeDesc range = new TFileRangeDesc();
+                    setIcebergParams(node, range, split);
+                    TIcebergFileDesc file = 
range.getTableFormatParams().getIcebergParams();
+                    Assert.assertTrue(file.isSetPartitionSpecId());
+                    Assert.assertEquals(task.file().specId(), 
file.getPartitionSpecId());
+                    Assert.assertFalse(file.isSetPartitionDataJson());
+                    if (keepPartitioned) {
+                        
Assert.assertEquals(Collections.singletonMap("record_key",
+                                        task.file().partition().get(0, 
Integer.class).toString()),
+                                split.getIcebergPartitionValues());
+                    }
+                    fileCount++;
+                }
+                Assert.assertEquals(2, fileCount);
+            }
+        }
+    }
+
+    @Test
+    public void testReadSplitAfterDroppingPartitionSourceColumn() throws 
Exception {
+        for (boolean keepPartitioned : Arrays.asList(false, true)) {
+            for (Integer value : Arrays.asList(7, null)) {
+                assertSplitAfterDroppingPartitionSourceColumn(keepPartitioned, 
value, false);
+            }
+        }
+    }
+
+    @Test
+    public void testRowIdScanRejectsDroppedPartitionSourceColumn() throws 
Exception {
+        for (boolean keepPartitioned : Arrays.asList(false, true)) {
+            for (Integer value : Arrays.asList(7, null)) {
+                assertSplitAfterDroppingPartitionSourceColumn(keepPartitioned, 
value, true);
+            }
+        }
+    }
+
+    private void assertSplitAfterDroppingPartitionSourceColumn(
+            boolean keepPartitioned, Integer value, boolean requireRowId) 
throws Exception {
+        assertSplitAfterDroppingPartitionSourceColumn(
+                keepPartitioned, value, requireRowId, new 
ExecutionAuthenticator() {});
+    }
+
+    @Test
+    public void testRowIdScanPreservesErrorThroughHadoopAuthentication() 
throws Exception {
+        assertSplitAfterDroppingPartitionSourceColumn(false, 7, true,
+                new HadoopExecutionAuthenticator(() -> 
UserGroupInformation.createRemoteUser("test_user")));
+    }
+
+    private void assertSplitAfterDroppingPartitionSourceColumn(boolean 
keepPartitioned, Integer value,
+            boolean requireRowId, ExecutionAuthenticator authenticator) throws 
Exception {
+        Schema schema = new Schema(Types.NestedField.optional(1, "record_key", 
Types.IntegerType.get()),
+                Types.NestedField.optional(2, "partition_key", 
Types.IntegerType.get()),
+                Types.NestedField.optional(3, "retained_key", 
Types.IntegerType.get()));
+        PartitionSpec.Builder specBuilder = 
PartitionSpec.builderFor(schema).identity("partition_key");
+        if (keepPartitioned) {
+            specBuilder.identity("retained_key");
+        }
+        HadoopTables tables = new HadoopTables(new Configuration());
+        String location = temporaryFolder.newFolder().toURI().toString();
+        Table table = tables.create(schema, specBuilder.build(),
+                Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"), 
location);
+        PartitionData partition = new 
PartitionData(table.spec().partitionType());
+        partition.set(0, value);
+        if (keepPartitioned) {
+            partition.set(1, 3);
+        }
+        
table.newAppend().appendFile(DataFiles.builder(table.spec()).withPath("file:///warehouse/old.parquet")
+                
.withPartition(partition).withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+        table.updateSpec().removeField("partition_key").commit();
+        table.updateSchema().deleteColumn("partition_key").commit();
+        // Reload metadata so historical specs bind to the schema without the 
dropped source column.
+        table = tables.load(location);
+        Assert.assertEquals(keepPartitioned, table.spec().isPartitioned());
+        TestIcebergScanNode node = new TestIcebergScanNode(new 
SessionVariable());
+        setIcebergTable(node, table);
+        setPrivateField(node, "isPartitionedTable", keepPartitioned);
+        setPrivateField(node, "storagePropertiesMap", Collections.emptyMap());
+        setPrivateField(node, "formatVersion", 2);
+        setPrivateField(node, "partitionMapInfos", new HashMap<>());
+        setPrivateField(node, "orderedPathPartitionKeys", 
Collections.emptyList());
+        setPrivateField(node, "orderedPartitionMetadataKeys", 
Collections.emptyList());
+        node.addSlot(0, new Column("record_key", Type.INT));
+        if (requireRowId) {
+            node.addSlot(1, IcebergRowId.createHiddenColumn());
+        }
+        try (CloseableIterable<FileScanTask> tasks = 
table.newScan().planFiles()) {
+            int fileCount = 0;
+            for (FileScanTask task : tasks) {
+                PartitionData actualPartition = (PartitionData) 
task.file().partition();
+                Assert.assertEquals(Types.UnknownType.get(), 
actualPartition.getPartitionType().fields().get(0).type());
+                Assert.assertEquals(value, actualPartition.get(0));
+                if (requireRowId) {
+                    UserException thrown = 
Assert.assertThrows(UserException.class,
+                            () -> getSplitsWithTask(node, task, 
authenticator));
+                    Assert.assertTrue(thrown.getMessage().contains(
+                            "Cannot produce Iceberg row IDs with unsupported 
partition types"));
+                    Assert.assertTrue(thrown.getMessage().contains("spec " + 
task.file().specId()));
+                    Assert.assertTrue(thrown.getMessage().contains("Rewrite 
historical data files"));
+                    Assert.assertTrue(thrown.getCause() instanceof 
UnsupportedOperationException);
+                } else {
+                    IcebergSplit split = createIcebergSplit(node, task);
+                    TFileRangeDesc range = new TFileRangeDesc();
+                    setIcebergParams(node, range, split);
+                    TIcebergFileDesc file = 
range.getTableFormatParams().getIcebergParams();
+                    Assert.assertEquals(task.file().specId(), 
file.getPartitionSpecId());
+                    Assert.assertFalse(file.isSetPartitionDataJson());
+                    if (keepPartitioned) {
+                        
Assert.assertEquals(Collections.singletonMap("retained_key", "3"),
+                                split.getIcebergPartitionValues());
+                    }
+                }
+                fileCount++;
+            }
+            Assert.assertEquals(1, fileCount);
+        }
+    }
+
+    private static void getSplitsWithTask(IcebergScanNode node, FileScanTask 
task,
+            ExecutionAuthenticator authenticator) throws Exception {
+        ConnectContext previous = ConnectContext.get();
+        ConnectContext context = new ConnectContext();
+        context.setStatementContext(new StatementContext());
+        
context.getStatementContext().setIcebergRewriteFileScanTasks(Collections.singletonList(task));
+        context.setThreadLocalInfo();
+        setPreExecutionAuthenticator(node, authenticator);
+        try {
+            // Exercise the public non-batch entry point used by row-level 
DML, including error wrapping.
+            node.getSplits(1);
+        } finally {
+            ConnectContext.remove();
+            if (previous != null) {
+                previous.setThreadLocalInfo();
+            }
+        }
+    }
+
+    @Test
+    public void testDeleteSplitKeepsBinaryPartitionAfterDroppingField() throws 
Exception {
+        assertDeleteAfterDroppingPartitionField(Types.BinaryType.get(),
+                ByteBuffer.wrap(new byte[] {0, (byte) 0xff, (byte) 0x80, 
0x2f}));
+        assertDeleteAfterDroppingPartitionField(Types.BinaryType.get(), 
ByteBuffer.allocate(0));
+        assertDeleteAfterDroppingPartitionField(Types.BinaryType.get(), null);
+    }
+
+    @Test
+    public void testDeleteSplitKeepsFixedPartitionAfterDroppingField() throws 
Exception {
+        assertDeleteAfterDroppingPartitionField(Types.FixedType.ofLength(4),
+                ByteBuffer.wrap(new byte[] {0, (byte) 0xff, (byte) 0x80, 
0x2f}));
+        assertDeleteAfterDroppingPartitionField(Types.FixedType.ofLength(4), 
null);
+    }
+
+    @Test
+    public void testDeleteSplitKeepsUuidPartitionAfterDroppingField() throws 
Exception {
+        assertDeleteAfterDroppingPartitionField(Types.UUIDType.get(),
+                UUID.fromString("123e4567-e89b-12d3-a456-426614174000"));
+    }
+
+    @Test
+    public void testDeleteSplitKeepsTimePartitionAfterDroppingField() throws 
Exception {
+        assertDeleteAfterDroppingPartitionField(Types.TimeType.get(), 
12_345_678_901L);
+    }
+
+    @Test
+    public void testDeleteSplitKeepsPreEpochTimestampAfterDroppingField() 
throws Exception {
+        
assertDeleteAfterDroppingPartitionField(Types.TimestampType.withoutZone(), -1L);
+        
assertDeleteAfterDroppingPartitionField(Types.TimestampType.withoutZone(), 
-1_000_001L);
+    }
+
+    @Test
+    public void testDeleteSplitKeepsPreEpochTimestamptzAfterDroppingField() 
throws Exception {
+        
assertDeleteAfterDroppingPartitionField(Types.TimestampType.withZone(), -1L);
+        
assertDeleteAfterDroppingPartitionField(Types.TimestampType.withZone(), 
-1_000_001L);
+    }
+
+    private void 
assertDeleteAfterDroppingPartitionField(org.apache.iceberg.types.Type type, 
Object value)
+            throws Exception {
+        Schema schema = new Schema(Types.NestedField.optional(1, 
"partition_key", type));
+        Table table = new HadoopTables(new Configuration()).create(schema,
+                
PartitionSpec.builderFor(schema).identity("partition_key").build(),
+                Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"),
+                temporaryFolder.newFolder().toURI().toString());
+        PartitionData partition = new 
PartitionData(table.spec().partitionType());
+        partition.set(0, value);
+        
table.newAppend().appendFile(DataFiles.builder(table.spec()).withPath("file:///warehouse/old.parquet")
+                
.withPartition(partition).withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+        table.updateSpec().removeField("partition_key").commit();
+        
table.newAppend().appendFile(DataFiles.builder(table.spec()).withPath("file:///warehouse/new.parquet")
+                .withRecordCount(1).withFileSizeInBytes(128).build()).commit();
+
+        Assert.assertTrue(table.spec().isUnpartitioned());
+        assertDeleteSplitPartitionMetadata(table);
+    }
+
+    private void assertDeleteSplitPartitionMetadata(Table table) throws 
Exception {
+        ConnectContext previousContext = ConnectContext.get();
+        ConnectContext context = new ConnectContext();
+        context.getSessionVariable().setTimeZone("UTC");
+        context.setThreadLocalInfo();
+        try {
+            assertDeleteSplitPartitionMetadata(table, 
context.getSessionVariable());
+        } finally {
+            if (previousContext == null) {
+                ConnectContext.remove();
+            } else {
+                previousContext.setThreadLocalInfo();
+            }
+        }
+    }
+
+    private void assertDeleteSplitPartitionMetadata(Table table, 
SessionVariable sessionVariable) throws Exception {
+        TestIcebergScanNode node = new TestIcebergScanNode(sessionVariable);
+        node.addSlot(0, IcebergRowId.createHiddenColumn());
+        setIcebergTable(node, table);
+        setPrivateField(node, "isPartitionedTable", 
table.spec().isPartitioned());
+        setPrivateField(node, "storagePropertiesMap", Collections.emptyMap());
+        setPrivateField(node, "formatVersion", 2);
+        setPrivateField(node, "partitionMapInfos", new HashMap<>());
+        setPrivateField(node, "orderedPathPartitionKeys", 
Collections.emptyList());
+        setPrivateField(node, "orderedPartitionMetadataKeys", 
Collections.emptyList());
+        try (CloseableIterable<FileScanTask> tasks = 
table.newScan().planFiles()) {
+            int fileCount = 0;
+            for (FileScanTask task : tasks) {
+                IcebergSplit split = createIcebergSplit(node, task);
+                TFileRangeDesc range = new TFileRangeDesc();
+                setIcebergParams(node, range, split);
+                byte[] bytes = new TSerializer(new 
TCompactProtocol.Factory()).serialize(range);
+                TFileRangeDesc restored = new TFileRangeDesc();
+                new TDeserializer(new 
TCompactProtocol.Factory()).deserialize(restored, bytes);
+                TIcebergFileDesc file = 
restored.getTableFormatParams().getIcebergParams();
+                Assert.assertTrue("Every data file must carry its spec id 
through Thrift", file.isSetPartitionSpecId());
+                Assert.assertEquals(task.file().specId(), 
file.getPartitionSpecId());
+                PartitionSpec spec = 
table.specs().get(file.getPartitionSpecId());
+                Assert.assertNotNull(file.getPartitionDataJson());
+                if (spec.isUnpartitioned()) {
+                    Assert.assertEquals("[]", file.getPartitionDataJson());
+                }
+
+                TIcebergCommitData commit = new TIcebergCommitData();
+                commit.setFilePath("delete-" + fileCount + ".parquet");
+                commit.setFileSize(128);
+                commit.setRowCount(1);
+                commit.setFileContent(TFileContent.POSITION_DELETES);
+                commit.setPartitionSpecId(file.getPartitionSpecId());
+                commit.setPartitionDataJson(file.getPartitionDataJson());
+                DeleteFile delete = IcebergWriterHelper.convertToDeleteFiles(
+                        FileFormat.PARQUET, spec, 
Collections.singletonList(commit)).get(0);
+                Assert.assertEquals(task.file().specId(), delete.specId());
+                Assert.assertEquals(task.file().partition().size(), 
delete.partition().size());
+                // PartitionData.equals compares binary backing arrays by 
identity, not byte content.
+                for (int i = 0; i < task.file().partition().size(); i++) {
+                    Assert.assertEquals(task.file().partition().get(i, 
Object.class),
+                            delete.partition().get(i, Object.class));
+                }
+                fileCount++;
+            }
+            
Assert.assertEquals(table.currentSnapshot().summary().get("total-data-files"),
+                    Integer.toString(fileCount));
+        }
+    }
+
     @Test
     public void testDeleteFileSizePropagatedToThrift() throws Exception {
         Types.NestedField id = Types.NestedField.required(1, "id", 
Types.LongType.get());
diff --git 
a/regression-test/suites/external_table_p0/iceberg/write/test_iceberg_delete_unpartitioned_evolution.groovy
 
b/regression-test/suites/external_table_p0/iceberg/write/test_iceberg_delete_unpartitioned_evolution.groovy
new file mode 100644
index 00000000000..138b3a6c82c
--- /dev/null
+++ 
b/regression-test/suites/external_table_p0/iceberg/write/test_iceberg_delete_unpartitioned_evolution.groovy
@@ -0,0 +1,198 @@
+// 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.
+
+suite("test_iceberg_delete_unpartitioned_evolution",
+        "p0,external,iceberg,external_docker,external_docker_iceberg") {
+    if 
(!"true".equalsIgnoreCase(context.config.otherConfigs.get("enableIcebergTest")))
 {
+        logger.info("disable iceberg test")
+        return
+    }
+
+    String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+    String restPort = context.config.otherConfigs.get("iceberg_rest_uri_port")
+    String minioPort = context.config.otherConfigs.get("iceberg_minio_port")
+    String catalogName = "test_iceberg_delete_unpartitioned_evolution"
+    String dbName = "iceberg_delete_unpartitioned_evolution_db"
+    sql """drop catalog if exists ${catalogName}"""
+    sql """
+        create catalog ${catalogName} properties (
+            "type" = "iceberg",
+            "iceberg.catalog.type" = "rest",
+            "uri" = "http://${externalEnvIp}:${restPort}";,
+            "s3.access_key" = "admin",
+            "s3.secret_key" = "password",
+            "s3.endpoint" = "http://${externalEnvIp}:${minioPort}";,
+            "s3.region" = "us-east-1"
+        )
+    """
+    sql """switch ${catalogName}"""
+    sql """create database if not exists ${dbName}"""
+    sql """use ${dbName}"""
+
+    def checkRows = { String table, def expected ->
+        sql """refresh table ${dbName}.${table}"""
+        spark_iceberg """refresh table demo.${dbName}.${table}"""
+        def actual = sql """select record_key, metric from ${table} order by 
record_key"""
+        assertSparkDorisResultEquals(expected, actual)
+        assertSparkDorisResultEquals(spark_iceberg("""
+            select record_key, metric from demo.${dbName}.${table} order by 
record_key
+        """), actual)
+    }
+
+    // REPLACE retains historical spec 0 while writing new files under an 
unpartitioned spec.
+    // Spark's replace transaction exercises the same Iceberg metadata 
transition as Trino RTAS.
+    spark_iceberg """drop table if exists 
demo.${dbName}.replace_to_unpartitioned"""
+    spark_iceberg """
+        create table demo.${dbName}.replace_to_unpartitioned (record_key int, 
metric int)
+        using iceberg partitioned by (record_key)
+        tblproperties ('format-version' = '2', 'write.format.default' = 
'parquet')
+    """
+    sql """refresh database ${dbName}"""
+    sql """insert into replace_to_unpartitioned values (1, 10), (2, 20), (3, 
30)"""
+    spark_iceberg """refresh table demo.${dbName}.replace_to_unpartitioned"""
+    spark_iceberg """
+        create or replace table demo.${dbName}.replace_to_unpartitioned using 
iceberg
+        tblproperties ('format-version' = '2', 'write.delete.mode' = 
'merge-on-read')
+        as select * from demo.${dbName}.replace_to_unpartitioned
+    """
+    sql """refresh table ${dbName}.replace_to_unpartitioned"""
+    def replacedSpecs = sql """select distinct spec_id from 
replace_to_unpartitioned\$data_files"""
+    assertEquals(1, replacedSpecs.size())
+    assertTrue((replacedSpecs[0][0] as int) > 0)
+    def expectedAfterReplaceDelete = spark_iceberg """
+        select record_key, metric from demo.${dbName}.replace_to_unpartitioned
+        where record_key <> 1 order by record_key
+    """
+    sql """delete from replace_to_unpartitioned where record_key = 1"""
+    checkRows("replace_to_unpartitioned", expectedAfterReplaceDelete)
+    assertEquals(replacedSpecs, sql("""
+        select distinct spec_id from replace_to_unpartitioned\$delete_files
+    """))
+    def expectedAfterRefreshDelete = spark_iceberg """
+        select record_key, metric from demo.${dbName}.replace_to_unpartitioned
+        where record_key <> 2 order by record_key
+    """
+    sql """delete from replace_to_unpartitioned where record_key = 2"""
+    checkRows("replace_to_unpartitioned", expectedAfterRefreshDelete)
+
+    // Dropping the last partition field leaves both old partitioned and new 
unpartitioned files live.
+    spark_iceberg """drop table if exists 
demo.${dbName}.mixed_partition_specs"""
+    spark_iceberg """
+        create table demo.${dbName}.mixed_partition_specs (record_key int, 
metric int)
+        using iceberg partitioned by (record_key)
+        tblproperties ('format-version' = '2', 'write.format.default' = 'orc',
+                       'write.delete.mode' = 'merge-on-read', 
'write.update.mode' = 'merge-on-read')
+    """
+    spark_iceberg """insert into demo.${dbName}.mixed_partition_specs values 
(1, 10), (2, 20)"""
+    spark_iceberg """alter table demo.${dbName}.mixed_partition_specs drop 
partition field record_key"""
+    spark_iceberg """insert into demo.${dbName}.mixed_partition_specs values 
(3, 30), (4, 40)"""
+    sql """refresh database ${dbName}"""
+    def mixedSpecs = sql """select distinct spec_id from 
mixed_partition_specs\$data_files order by spec_id"""
+    assertEquals(2, mixedSpecs.size())
+    def expectedAfterMixedDelete = spark_iceberg """
+        select record_key, metric from demo.${dbName}.mixed_partition_specs
+        where record_key not in (1, 3) order by record_key
+    """
+    sql """delete from mixed_partition_specs where record_key in (1, 3)"""
+    checkRows("mixed_partition_specs", expectedAfterMixedDelete)
+    assertEquals(mixedSpecs, sql("""
+        select distinct spec_id from mixed_partition_specs\$delete_files order 
by spec_id
+    """))
+    def expectedAfterUpdate = spark_iceberg """
+        select record_key, metric + 100 from 
demo.${dbName}.mixed_partition_specs order by record_key
+    """
+    sql """update mixed_partition_specs set metric = metric + 100"""
+    checkRows("mixed_partition_specs", expectedAfterUpdate)
+
+    // Dropping a source column makes its historical partition type unknown, 
even for non-null values.
+    for (String fileFormat : ["parquet", "orc"]) {
+        String table = "dropped_partition_source_${fileFormat}"
+        spark_iceberg """drop table if exists demo.${dbName}.${table}"""
+        spark_iceberg """
+            create table demo.${dbName}.${table} (record_key int, 
partition_key int, metric int)
+            using iceberg partitioned by (partition_key)
+            tblproperties ('format-version' = '2', 'write.format.default' = 
'${fileFormat}',
+                           'write.delete.mode' = 'merge-on-read', 
'write.update.mode' = 'merge-on-read')
+        """
+        spark_iceberg """insert into demo.${dbName}.${table} values (1, 7, 
10), (2, 7, 20), (3, null, 30)"""
+        spark_iceberg """alter table demo.${dbName}.${table} drop partition 
field partition_key"""
+        spark_iceberg """alter table demo.${dbName}.${table} drop column 
partition_key"""
+        spark_iceberg """insert into demo.${dbName}.${table} values (4, 40)"""
+        sql """refresh database ${dbName}"""
+        def expected = spark_iceberg """
+            select record_key, metric from demo.${dbName}.${table} order by 
record_key
+        """
+        checkRows(table, expected)
+        test {
+            sql """delete from ${table} where record_key = 1"""
+            exception "Cannot produce Iceberg row IDs with unsupported 
partition types"
+        }
+        test {
+            sql """update ${table} set metric = metric + 100 where record_key 
= 2"""
+            exception "Cannot produce Iceberg row IDs with unsupported 
partition types"
+        }
+        checkRows(table, expected)
+        assertEquals([[0L]], sql("""select count(*) from 
${table}\$delete_files"""))
+    }
+
+    // Historical binary partitions must retain their bytes, not become NULL 
delete partitions.
+    spark_iceberg """drop table if exists demo.${dbName}.historical_binary"""
+    spark_iceberg """
+        create table demo.${dbName}.historical_binary (record_key int, 
partition_key binary, metric int)
+        using iceberg partitioned by (partition_key)
+        tblproperties ('format-version' = '2', 'write.delete.mode' = 
'merge-on-read')
+    """
+    spark_iceberg """insert into demo.${dbName}.historical_binary values
+        (1, unhex('00ff802f'), 10), (2, unhex('00ff802f'), 20), (3, null, 30), 
(4, unhex(''), 40)"""
+    spark_iceberg """alter table demo.${dbName}.historical_binary drop 
partition field partition_key"""
+    spark_iceberg """insert into demo.${dbName}.historical_binary values (5, 
unhex('ff'), 50)"""
+    sql """refresh database ${dbName}"""
+    def expectedAfterBinaryDelete = spark_iceberg """
+        select record_key, metric from demo.${dbName}.historical_binary
+        where record_key not in (1, 3, 4, 5) order by record_key
+    """
+    sql """delete from historical_binary where record_key in (1, 3, 4, 5)"""
+    checkRows("historical_binary", expectedAfterBinaryDelete)
+
+    // Negative fractional epochs require floor-based splitting into seconds 
and nanoseconds.
+    spark_iceberg """drop table if exists 
demo.${dbName}.historical_timestamp"""
+    spark_iceberg """
+        create table demo.${dbName}.historical_timestamp
+            (record_key int, partition_key timestamp_ntz, metric int)
+        using iceberg partitioned by (partition_key)
+        tblproperties ('format-version' = '2', 'write.delete.mode' = 
'merge-on-read',
+                       'write.update.mode' = 'merge-on-read')
+    """
+    spark_iceberg """insert into demo.${dbName}.historical_timestamp values
+        (1, cast('1969-12-31 23:59:59.999999' as timestamp_ntz), 10),
+        (2, cast('1969-12-31 23:59:58.999999' as timestamp_ntz), 20), (3, 
null, 30)"""
+    spark_iceberg """alter table demo.${dbName}.historical_timestamp drop 
partition field partition_key"""
+    spark_iceberg """insert into demo.${dbName}.historical_timestamp values
+        (4, cast('1970-01-01 00:00:00.000001' as timestamp_ntz), 40)"""
+    sql """refresh database ${dbName}"""
+    def expectedAfterTimestampDelete = spark_iceberg """
+        select record_key, metric from demo.${dbName}.historical_timestamp
+        where record_key not in (1, 4) order by record_key
+    """
+    sql """delete from historical_timestamp where record_key in (1, 4)"""
+    checkRows("historical_timestamp", expectedAfterTimestampDelete)
+    def expectedAfterTimestampUpdate = spark_iceberg """
+        select record_key, metric + 100 from 
demo.${dbName}.historical_timestamp order by record_key
+    """
+    sql """update historical_timestamp set metric = metric + 100"""
+    checkRows("historical_timestamp", expectedAfterTimestampUpdate)
+}


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

Reply via email to