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]