This is an automated email from the ASF dual-hosted git repository.

ahmedabu98 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new f5d35728554 [IcebergIO] Spec aware file reserialization (#39973)
f5d35728554 is described below

commit f5d3572855405b65208f0f163abbdde9cb6a6278
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Thu Sep 3 20:25:29 2026 +0200

    [IcebergIO] Spec aware file reserialization (#39973)
    
    * spec-aware
    
    * fix(iceberg): drop the Logger orphaned by the fallback-warning removal
    
    ErrorProne -Werror fails compileJava on the unused field, and the test
    javadoc linked the deleted method.
    
    Co-Authored-By: Claude Fable 5 <[email protected]>
    
    ---------
    
    Co-authored-by: Ahmed Abualsaud <[email protected]>
    Co-authored-by: Claude Fable 5 <[email protected]>
---
 .../sdk/io/iceberg/SerializableDeleteFile.java     |  85 ++++++-
 .../io/iceberg/cdc/SerializableChangelogTask.java  |  10 +-
 .../sdk/io/iceberg/SerializableDeleteFileTest.java | 281 +++++++++++++++++++--
 .../beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java  |   5 +-
 .../sdk/io/iceberg/cdc/ReadFromChangelogsTest.java |   5 +-
 5 files changed, 347 insertions(+), 39 deletions(-)

diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFile.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFile.java
index ceb96d50f8a..779bc3753c7 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFile.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFile.java
@@ -28,6 +28,7 @@ import java.nio.ByteBuffer;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
+import org.apache.beam.sdk.annotations.Internal;
 import org.apache.beam.sdk.schemas.AutoValueSchema;
 import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
 import org.apache.beam.sdk.schemas.annotations.SchemaFieldNumber;
@@ -37,11 +38,14 @@ import org.apache.iceberg.FileFormat;
 import org.apache.iceberg.FileMetadata;
 import org.apache.iceberg.Metrics;
 import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.SingleValueParser;
 import org.apache.iceberg.SortOrder;
+import org.apache.iceberg.StructLike;
 import org.checkerframework.checker.nullness.qual.Nullable;
 
 @DefaultSchema(AutoValueSchema.class)
 @AutoValue
+@Internal
 public abstract class SerializableDeleteFile {
   public static SerializableDeleteFile.Builder builder() {
     return new AutoValue_SerializableDeleteFile.Builder();
@@ -62,7 +66,11 @@ public abstract class SerializableDeleteFile {
   @SchemaFieldNumber("4")
   public abstract long getFileSizeInBytes();
 
+  /**
+   * @deprecated Use {@link #getJsonPartition()} instead.
+   */
   @SchemaFieldNumber("5")
+  @Deprecated
   public abstract String getPartitionPath();
 
   @SchemaFieldNumber("6")
@@ -113,6 +121,9 @@ public abstract class SerializableDeleteFile {
   @SchemaFieldNumber("21")
   public abstract @Nullable Long getFileSequenceNumber();
 
+  @SchemaFieldNumber("22")
+  abstract @Nullable String getJsonPartition();
+
   @AutoValue.Builder
   abstract static class Builder {
     abstract Builder setContentType(FileContent content);
@@ -127,6 +138,8 @@ public abstract class SerializableDeleteFile {
 
     abstract Builder setPartitionPath(String partitionPath);
 
+    abstract Builder setJsonPartition(String jsonPartition);
+
     abstract Builder setPartitionSpecId(int partitionSpec);
 
     abstract Builder setSortOrderId(@Nullable Integer sortOrderId);
@@ -163,7 +176,47 @@ public abstract class SerializableDeleteFile {
   }
 
   public static SerializableDeleteFile from(
-      DeleteFile deleteFile, String partitionPath, boolean includeMetrics) {
+      DeleteFile deleteFile, Map<Integer, PartitionSpec> specs) {
+    return from(deleteFile, specs, true);
+  }
+
+  /**
+   * Creates a {@link SerializableDeleteFile}, resolving the file's {@link 
PartitionSpec} by its own
+   * spec id.
+   *
+   * <p>Delete files reached from a scan task may carry a spec id that differs 
from the spec of the
+   * data file they apply to, so the lookup has to be per delete file rather 
than against a single
+   * "current" spec.
+   */
+  public static SerializableDeleteFile from(
+      DeleteFile deleteFile, Map<Integer, PartitionSpec> specs, boolean 
includeMetrics) {
+    return from(
+        deleteFile,
+        checkStateNotNull(
+            specs.get(deleteFile.specId()),
+            "Could not create a SerializableDeleteFile because DeleteFile is 
written using a partition spec id '%s' that is not found in the provided specs: 
%s",
+            deleteFile.specId(),
+            specs.keySet()),
+        includeMetrics);
+  }
+
+  public static SerializableDeleteFile from(DeleteFile deleteFile, 
PartitionSpec spec) {
+    return from(deleteFile, spec, true);
+  }
+
+  public static SerializableDeleteFile from(
+      DeleteFile deleteFile, PartitionSpec spec, boolean includeMetrics) {
+    if (spec.specId() != deleteFile.specId()) {
+      throw new IllegalArgumentException(
+          String.format(
+              "Cannot serialize DeleteFile: its partition spec id %s does not 
match the provided "
+                  + "spec id %s.",
+              deleteFile.specId(), spec.specId()));
+    }
+    // jsonPartition is the primary (handles evolved specs, special 
characters).
+    // partitionPath is the fallback for values that don't round-trip through 
JSON.
+    String jsonPartition = SingleValueParser.toJson(spec.partitionType(), 
deleteFile.partition());
+    String partitionPath = spec.partitionToPath(deleteFile.partition());
 
     SerializableDeleteFile.Builder builder =
         SerializableDeleteFile.builder()
@@ -171,6 +224,7 @@ public abstract class SerializableDeleteFile {
             .setFileFormat(deleteFile.format().name())
             .setFileSizeInBytes(deleteFile.fileSizeInBytes())
             .setPartitionPath(partitionPath)
+            .setJsonPartition(jsonPartition)
             .setPartitionSpecId(deleteFile.specId())
             .setRecordCount(deleteFile.recordCount())
             .setColumnSizes(deleteFile.columnSizes())
@@ -228,7 +282,21 @@ public abstract class SerializableDeleteFile {
             .withMetrics(metrics)
             .withSplitOffsets(getSplitOffsets())
             .withEncryptionKeyMetadata(getKeyMetadata())
-            .withPartitionPath(getPartitionPath());
+            .withReferencedDataFile(getReferencedDataFile());
+
+    @Nullable String jsonPartition = getJsonPartition();
+    if (jsonPartition != null) {
+      try {
+        deleteFileBuilder = 
deleteFileBuilder.withPartition(partition(partitionSpec));
+      } catch (RuntimeException e) {
+        // Some partition values (e.g. NaN / Infinity floating-point) don't 
round-trip through the
+        // JSON representation; fall back to the partition-path string
+        deleteFileBuilder = 
deleteFileBuilder.withPartitionPath(getPartitionPath());
+      }
+    } else {
+      // Elements decoded from a pre-jsonPartition release carry only the 
partition path.
+      deleteFileBuilder = 
deleteFileBuilder.withPartitionPath(getPartitionPath());
+    }
 
     switch (getContentType()) {
       case POSITION_DELETES:
@@ -260,17 +328,22 @@ public abstract class SerializableDeleteFile {
             "Unexpected content type for DeleteFile: " + getContentType());
     }
 
-    // needed for puffin files
+    // contentOffset / contentSizeInBytes really are Puffin-only: build() 
rejects a non-null value
+    // for either on any other format, and requires both (plus 
referencedDataFile) on Puffin.
     if (getFileFormat().equalsIgnoreCase(FileFormat.PUFFIN.name())) {
       deleteFileBuilder =
           deleteFileBuilder
               .withContentOffset(checkStateNotNull(getContentOffset()))
-              
.withContentSizeInBytes(checkStateNotNull(getContentSizeInBytes()))
-              
.withReferencedDataFile(checkStateNotNull(getReferencedDataFile()));
+              
.withContentSizeInBytes(checkStateNotNull(getContentSizeInBytes()));
     }
     return deleteFileBuilder.build();
   }
 
+  private StructLike partition(PartitionSpec spec) {
+    return (StructLike)
+        SingleValueParser.fromJson(spec.partitionType(), 
checkStateNotNull(getJsonPartition()));
+  }
+
   @Override
   public final boolean equals(@Nullable Object o) {
     if (this == o) {
@@ -287,6 +360,7 @@ public abstract class SerializableDeleteFile {
         && getFileSizeInBytes() == that.getFileSizeInBytes()
         && getPartitionPath().equals(that.getPartitionPath())
         && getPartitionSpecId() == that.getPartitionSpecId()
+        && Objects.equals(getJsonPartition(), that.getJsonPartition())
         && Objects.equals(getSortOrderId(), that.getSortOrderId())
         && Objects.equals(getEqualityFieldIds(), that.getEqualityFieldIds())
         && Objects.equals(getKeyMetadata(), that.getKeyMetadata())
@@ -314,6 +388,7 @@ public abstract class SerializableDeleteFile {
             getRecordCount(),
             getFileSizeInBytes(),
             getPartitionPath(),
+            getJsonPartition(),
             getPartitionSpecId(),
             getSortOrderId(),
             getEqualityFieldIds(),
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
index 3410c0a9d7e..97bcbaa5bec 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
@@ -17,7 +17,6 @@
  */
 package org.apache.beam.sdk.io.iceberg.cdc;
 
-import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState;
 
 import com.google.auto.value.AutoValue;
@@ -267,13 +266,10 @@ public abstract class SerializableChangelogTask {
 
   private static List<SerializableDeleteFile> toSerializableDeletes(
       List<DeleteFile> dfs, Map<Integer, PartitionSpec> specs, boolean 
includeMetrics) {
+    // Serialize each delete file against its own spec (looked up by its spec 
id): a delete file may
+    // carry a different spec id than the data file it applies to.
     return dfs.stream()
-        .map(
-            df ->
-                SerializableDeleteFile.from(
-                    df,
-                    
checkStateNotNull(specs.get(df.specId())).partitionToPath(df.partition()),
-                    includeMetrics))
+        .map(df -> SerializableDeleteFile.from(df, specs, includeMetrics))
         .collect(Collectors.toList());
   }
 }
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFileTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFileTest.java
index 29ef30c97ef..02e6fec4657 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFileTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableDeleteFileTest.java
@@ -19,6 +19,12 @@ package org.apache.beam.sdk.io.iceberg;
 
 import static java.util.Collections.emptyMap;
 import static java.util.Collections.singletonMap;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsString;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.instanceOf;
+import static org.hamcrest.Matchers.not;
+import static org.hamcrest.Matchers.nullValue;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNull;
 import static org.junit.Assert.assertThrows;
@@ -30,6 +36,9 @@ import java.nio.ByteBuffer;
 import java.util.Arrays;
 import java.util.HashMap;
 import java.util.Map;
+import org.apache.beam.sdk.schemas.SchemaCoder;
+import org.apache.beam.sdk.schemas.SchemaRegistry;
+import org.apache.beam.sdk.util.CoderUtils;
 import org.apache.iceberg.DeleteFile;
 import org.apache.iceberg.FileContent;
 import org.apache.iceberg.FileFormat;
@@ -37,10 +46,16 @@ import org.apache.iceberg.FileMetadata;
 import org.apache.iceberg.Metrics;
 import org.apache.iceberg.PartitionSpec;
 import org.apache.iceberg.SortOrder;
+import org.apache.iceberg.StructLike;
+import org.apache.iceberg.data.GenericRecord;
 import org.apache.iceberg.types.Types;
+import org.checkerframework.checker.nullness.qual.Nullable;
 import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
 
 /** Tests for {@link SerializableDeleteFile}. */
+@RunWith(JUnit4.class)
 public class SerializableDeleteFileTest {
   private static final org.apache.iceberg.Schema SCHEMA =
       new org.apache.iceberg.Schema(
@@ -85,7 +100,7 @@ public class SerializableDeleteFileTest {
             .build();
     setSequenceNumbers(deleteFile, 44L, 45L);
 
-    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, "category=A", true);
+    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, SPEC, true);
     DeleteFile reconstructed =
         serialized.createDeleteFile(
             singletonMap(SPEC.specId(), SPEC), singletonMap(0, 
SortOrder.unsorted()));
@@ -125,7 +140,7 @@ public class SerializableDeleteFileTest {
             .withRecordCount(2L)
             .build();
 
-    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, "category=A", true);
+    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, SPEC, true);
     DeleteFile reconstructed =
         serialized.createDeleteFile(singletonMap(SPEC.specId(), SPEC), 
singletonMap(7, sortOrder));
 
@@ -149,7 +164,7 @@ public class SerializableDeleteFileTest {
             .withReferencedDataFile("gs://bucket/data/category=A/data.parquet")
             .build();
 
-    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, "category=A", true);
+    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, SPEC, true);
     DeleteFile reconstructed =
         serialized.createDeleteFile(
             singletonMap(SPEC.specId(), SPEC), singletonMap(0, 
SortOrder.unsorted()));
@@ -160,9 +175,11 @@ public class SerializableDeleteFileTest {
     assertEquals("gs://bucket/data/category=A/data.parquet", 
reconstructed.referencedDataFile());
   }
 
+  /** Reconstruction fails clearly when the spec map or the sort-order map 
lacks the file's id. */
   @Test
-  public void testCreateDeleteFileFailsClearlyForMissingPartitionSpec() {
-    DeleteFile deleteFile =
+  public void testCreateDeleteFileFailsClearlyForMissingSpecOrSortOrder() {
+    // facet: missing partition spec.
+    DeleteFile positionDelete =
         FileMetadata.deleteFileBuilder(SPEC)
             .ofPositionDeletes()
             .withPath("gs://bucket/deletes/category=A/pos.parquet")
@@ -171,19 +188,17 @@ public class SerializableDeleteFileTest {
             .withFileSizeInBytes(256L)
             .withRecordCount(2L)
             .build();
-    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, "category=A", true);
-
-    IllegalStateException thrown =
+    SerializableDeleteFile serializedPosition =
+        SerializableDeleteFile.from(positionDelete, SPEC, true);
+    IllegalStateException missingSpec =
         assertThrows(
-            IllegalStateException.class, () -> 
serialized.createDeleteFile(emptyMap(), null));
-
-    assertTrue(thrown.getMessage().contains("created with spec id '" + 
SPEC.specId() + "'"));
-  }
+            IllegalStateException.class,
+            () -> serializedPosition.createDeleteFile(emptyMap(), null));
+    assertTrue(missingSpec.getMessage().contains("created with spec id '" + 
SPEC.specId() + "'"));
 
-  @Test
-  public void testCreateEqualityDeleteFileFailsClearlyForMissingSortOrder() {
+    // facet: missing sort order (equality delete).
     SortOrder sortOrder = 
SortOrder.builderFor(SCHEMA).asc("id").withOrderId(7).build();
-    DeleteFile deleteFile =
+    DeleteFile equalityDelete =
         FileMetadata.deleteFileBuilder(SPEC)
             .ofEqualityDeletes(1)
             .withSortOrder(sortOrder)
@@ -193,14 +208,242 @@ public class SerializableDeleteFileTest {
             .withFileSizeInBytes(256L)
             .withRecordCount(2L)
             .build();
-    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, "category=A", true);
+    SerializableDeleteFile serializedEquality =
+        SerializableDeleteFile.from(equalityDelete, SPEC, true);
+    IllegalStateException missingOrder =
+        assertThrows(
+            IllegalStateException.class,
+            () ->
+                
serializedEquality.createDeleteFile(singletonMap(SPEC.specId(), SPEC), 
emptyMap()));
+    assertTrue(missingOrder.getMessage().contains("sort order id '7'"));
+  }
 
-    IllegalStateException thrown =
+  /**
+   * A {@link DeleteFile} must be serialized with the EXACT spec it was 
written with: a mismatched
+   * spec id and a spec id absent from the map each fail loudly.
+   */
+  @Test
+  public void fromRejectsMismatchedOrUnknownSpec() {
+    // facet: single-spec overload, wrong spec.
+    DeleteFile unpartitioned =
+        FileMetadata.deleteFileBuilder(PartitionSpec.unpartitioned())
+            .ofPositionDeletes()
+            .withPath("gs://bucket/deletes/pos.parquet")
+            .withFormat(FileFormat.PARQUET)
+            .withFileSizeInBytes(1L)
+            .withRecordCount(1L)
+            .build();
+    PartitionSpec otherSpec =
+        
PartitionSpec.builderFor(SCHEMA).identity("category").withSpecId(1).build();
+    IllegalArgumentException mismatch =
+        assertThrows(
+            IllegalArgumentException.class,
+            () -> SerializableDeleteFile.from(unpartitioned, otherSpec));
+    assertThat(mismatch.getMessage(), containsString("does not match"));
+
+    // facet: spec-map overload, id missing from the map.
+    DeleteFile deleteFile = positionDeletes(SPEC, partition(SPEC, "A"));
+    IllegalStateException unknown =
         assertThrows(
             IllegalStateException.class,
-            () -> serialized.createDeleteFile(singletonMap(SPEC.specId(), 
SPEC), emptyMap()));
+            () -> SerializableDeleteFile.from(deleteFile, emptyMap(), true));
+    assertThat(unknown.getMessage(), containsString("partition spec id '0'"));
+  }
+
+  /**
+   * All three delete-file kinds — equality, position, and a V3 deletion 
vector — round-trip through
+   * the schema coder with their partition tuples intact.
+   */
+  @Test
+  public void allDeleteFileKindsRoundTripThroughSchemaCoderWithPartition() 
throws Exception {
+    // facet: equality delete.
+    SortOrder sortOrder = 
SortOrder.builderFor(SCHEMA).asc("id").withOrderId(7).build();
+    DeleteFile equalityDelete =
+        FileMetadata.deleteFileBuilder(SPEC)
+            .ofEqualityDeletes(1, 2)
+            .withSortOrder(sortOrder)
+            .withPath("gs://bucket/deletes/category=A/eq.parquet")
+            .withFormat(FileFormat.PARQUET)
+            .withPartition(partition(SPEC, "A"))
+            .withFileSizeInBytes(256L)
+            .withRecordCount(2L)
+            .build();
+    DeleteFile equalityReconstructed =
+        encodeDecode(SerializableDeleteFile.from(equalityDelete, SPEC))
+            .createDeleteFile(singletonMap(SPEC.specId(), SPEC), 
singletonMap(7, sortOrder));
+    assertEquals(FileContent.EQUALITY_DELETES, 
equalityReconstructed.content());
+    assertEquals(equalityDelete.partition(), 
equalityReconstructed.partition());
+    assertEquals("category=A", 
SPEC.partitionToPath(equalityReconstructed.partition()));
+
+    // facet: position delete.
+    DeleteFile positionDelete = positionDeletes(SPEC, partition(SPEC, "A"));
+    DeleteFile positionReconstructed =
+        encodeDecode(SerializableDeleteFile.from(positionDelete, SPEC))
+            .createDeleteFile(singletonMap(SPEC.specId(), SPEC), null);
+    assertEquals(FileContent.POSITION_DELETES, 
positionReconstructed.content());
+    assertEquals(positionDelete.partition(), 
positionReconstructed.partition());
+    assertEquals("category=A", 
SPEC.partitionToPath(positionReconstructed.partition()));
+
+    // facet: V3 deletion vector (a Puffin blob with 
offset/size/referenced-data-file).
+    DeleteFile dv =
+        FileMetadata.deleteFileBuilder(SPEC)
+            .ofPositionDeletes()
+            .withPath("gs://bucket/deletes/category=A/dv.puffin")
+            .withFormat(FileFormat.PUFFIN)
+            .withPartition(partition(SPEC, "A"))
+            .withFileSizeInBytes(512L)
+            .withRecordCount(1L)
+            .withContentOffset(64L)
+            .withContentSizeInBytes(128L)
+            .withReferencedDataFile("gs://bucket/data/category=A/data.parquet")
+            .build();
+    DeleteFile dvReconstructed =
+        encodeDecode(SerializableDeleteFile.from(dv, SPEC))
+            .createDeleteFile(singletonMap(SPEC.specId(), SPEC), null);
+    assertEquals(FileFormat.PUFFIN, dvReconstructed.format());
+    assertEquals(Long.valueOf(64L), dvReconstructed.contentOffset());
+    assertEquals(Long.valueOf(128L), dvReconstructed.contentSizeInBytes());
+    assertEquals("gs://bucket/data/category=A/data.parquet", 
dvReconstructed.referencedDataFile());
+    assertEquals(dv.partition(), dvReconstructed.partition());
+    assertEquals("category=A", 
SPEC.partitionToPath(dvReconstructed.partition()));
+  }
+
+  /**
+   * An identity partition on a {@code timestamptz} column: {@link
+   * org.apache.iceberg.types.Conversions#fromPartitionString} has no case for 
TIMESTAMP at all, so
+   * the old partition-path round-trip blew up at reconstruct time. The JSON 
representation carries
+   * the raw micros and round-trips exactly.
+   */
+  @Test
+  public void 
timestampPartitionRoundTripsThroughJsonButNotThroughPartitionPath() {
+    org.apache.iceberg.Schema schema =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.IntegerType.get()),
+            Types.NestedField.required(2, "event_time", 
Types.TimestampType.withZone()));
+    PartitionSpec spec = 
PartitionSpec.builderFor(schema).identity("event_time").build();
+    long micros = 1_709_618_828_000_009L;
+    GenericRecord partition = GenericRecord.create(spec.partitionType());
+    partition.setField("event_time", micros);
+    DeleteFile deleteFile = positionDeletes(spec, partition);
+
+    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, spec);
+    DeleteFile reconstructed = 
serialized.createDeleteFile(singletonMap(spec.specId(), spec), null);
+
+    assertEquals(Long.valueOf(micros), reconstructed.partition().get(0, 
Long.class));
+    assertEquals(deleteFile.partition(), reconstructed.partition());
+
+    // The pre-change wire shape (partition path only) cannot reconstruct this 
partition at all.
+    SerializableDeleteFile legacy = withoutJsonPartition(serialized);
+    assertThrows(
+        UnsupportedOperationException.class,
+        () -> legacy.createDeleteFile(singletonMap(spec.specId(), spec), 
null));
+  }
+
+  /**
+   * {@code PartitionSpec.partitionToPath} URL-encodes each value but {@code 
DataFiles.fillFromPath}
+   * never decodes it, so a string partition containing {@code / }, {@code &} 
or {@code =} used to
+   * come back SILENTLY WRONG — the delete would be registered under a {@code 
(specId, partition)}
+   * that no data file lives in, and would simply never apply. The JSON 
representation is exact.
+   */
+  @Test
+  public void stringPartitionWithSpecialCharactersRoundTripsExactly() {
+    String value = "a/b c&d=e";
+    DeleteFile deleteFile = positionDeletes(SPEC, partition(SPEC, value));
+
+    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, SPEC);
+    DeleteFile reconstructed = 
serialized.createDeleteFile(singletonMap(SPEC.specId(), SPEC), null);
+
+    assertEquals(value, reconstructed.partition().get(0, 
CharSequence.class).toString());
+    assertEquals(deleteFile.partition(), reconstructed.partition());
+
+    // Prove the old representation was silently lossy rather than merely 
throwing.
+    DeleteFile viaLegacyPath =
+        
withoutJsonPartition(serialized).createDeleteFile(singletonMap(SPEC.specId(), 
SPEC), null);
+    assertThat(
+        viaLegacyPath.partition().get(0, CharSequence.class).toString(), 
not(equalTo(value)));
+  }
+
+  /**
+   * A null partition value is rendered as the literal text {@code null} in a 
partition path, which
+   * {@code fromPartitionString} hands back as the four-character string 
"null" for a string column
+   * — again silently wrong. JSON omits the field and it decodes back to a 
real null.
+   */
+  @Test
+  public void nullPartitionValueRoundTripsAsNull() {
+    DeleteFile deleteFile = positionDeletes(SPEC, partition(SPEC, null));
+
+    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, SPEC);
+    DeleteFile reconstructed = 
serialized.createDeleteFile(singletonMap(SPEC.specId(), SPEC), null);
+
+    assertThat(reconstructed.partition().get(0, CharSequence.class), 
nullValue());
+    assertEquals(deleteFile.partition(), reconstructed.partition());
+
+    // The old path turns the null into the literal string "null".
+    DeleteFile viaLegacyPath =
+        
withoutJsonPartition(serialized).createDeleteFile(singletonMap(SPEC.specId(), 
SPEC), null);
+    assertEquals("null", viaLegacyPath.partition().get(0, 
CharSequence.class).toString());
+  }
+
+  /**
+   * NaN / Infinity floating-point partition values don't round-trip through 
the JSON partition
+   * representation ({@code SingleValueParser.fromJson} rejects the 
non-standard {@code NaN} token).
+   * Reconstruct must fall back to the partition-path string, which handles 
them, rather than
+   * crash-looping the sink at commit time.
+   */
+  @Test
+  public void nanFloatPartitionReconstructsViaPathFallback() {
+    org.apache.iceberg.Schema schema =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "f", Types.FloatType.get()),
+            Types.NestedField.optional(2, "data", Types.StringType.get()));
+    PartitionSpec spec = 
PartitionSpec.builderFor(schema).identity("f").build();
+    GenericRecord partition = GenericRecord.create(spec.partitionType());
+    partition.setField("f", Float.NaN);
+    DeleteFile deleteFile = positionDeletes(spec, partition);
+
+    SerializableDeleteFile serialized = 
SerializableDeleteFile.from(deleteFile, spec);
+    // Must reconstruct without throwing (JSON decode of NaN fails -> 
partition-path fallback).
+    DeleteFile reconstructed = 
serialized.createDeleteFile(singletonMap(spec.specId(), spec), null);
+
+    Object value = reconstructed.partition().get(0, Object.class);
+    assertThat(value, instanceOf(Float.class));
+    assertTrue("partition value must round-trip as NaN", Float.isNaN((Float) 
value));
+  }
+
+  private static GenericRecord partition(PartitionSpec spec, @Nullable Object 
value) {
+    GenericRecord record = GenericRecord.create(spec.partitionType());
+    record.set(0, value);
+    return record;
+  }
+
+  private static DeleteFile positionDeletes(PartitionSpec spec, StructLike 
partition) {
+    return FileMetadata.deleteFileBuilder(spec)
+        .ofPositionDeletes()
+        .withPath("gs://bucket/deletes/pos.parquet")
+        .withFormat(FileFormat.PARQUET)
+        .withPartition(partition)
+        .withFileSizeInBytes(256L)
+        .withRecordCount(2L)
+        .build();
+  }
+
+  /** Rebuilds the element as a pre-jsonPartition release would have encoded 
it. */
+  private static SerializableDeleteFile 
withoutJsonPartition(SerializableDeleteFile file) {
+    return SerializableDeleteFile.builder()
+        .setContentType(file.getContentType())
+        .setLocation(file.getLocation())
+        .setFileFormat(file.getFileFormat())
+        .setRecordCount(file.getRecordCount())
+        .setFileSizeInBytes(file.getFileSizeInBytes())
+        .setPartitionPath(file.getPartitionPath())
+        .setPartitionSpecId(file.getPartitionSpecId())
+        .build();
+  }
 
-    assertTrue(thrown.getMessage().contains("sort order id '7'"));
+  private static SerializableDeleteFile encodeDecode(SerializableDeleteFile 
file) throws Exception {
+    SchemaCoder<SerializableDeleteFile> coder =
+        
SchemaRegistry.createDefault().getSchemaCoder(SerializableDeleteFile.class);
+    return CoderUtils.decodeFromByteArray(coder, 
CoderUtils.encodeToByteArray(coder, file));
   }
 
   private static void setSequenceNumbers(
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java
index 75bfaa1775d..8c23a844be6 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java
@@ -360,10 +360,7 @@ public class CdcReadUtilsTest {
   private static List<SerializableDeleteFile> serializableDeletes(
       List<DeleteFile> deletes, Table table) {
     return deletes.stream()
-        .map(
-            delete ->
-                SerializableDeleteFile.from(
-                    delete, table.spec().partitionToPath(delete.partition()), 
true))
+        .map(delete -> SerializableDeleteFile.from(delete, table.specs(), 
true))
         .collect(Collectors.toList());
   }
 
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogsTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogsTest.java
index 7a1e71d3d80..2de2848f936 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogsTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogsTest.java
@@ -316,10 +316,7 @@ public class ReadFromChangelogsTest {
   private static List<SerializableDeleteFile> serializableDeletes(
       List<DeleteFile> deletes, Table table) {
     return deletes.stream()
-        .map(
-            delete ->
-                SerializableDeleteFile.from(
-                    delete, table.spec().partitionToPath(delete.partition()), 
true))
+        .map(delete -> SerializableDeleteFile.from(delete, table.specs(), 
true))
         .collect(Collectors.toList());
   }
 

Reply via email to