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