stevenzwu commented on code in PR #18248:
URL: https://github.com/apache/iceberg/pull/18248#discussion_r4097931939


##########
core/src/main/java/org/apache/iceberg/V4ManifestReader.java:
##########
@@ -113,11 +115,20 @@ Schema readSchema() {
 
   @Override
   public CloseableIterator<TrackedFile> iterator() {
+    // first row ID assignment must happen first to maintain consistent 
assignment when files are
+    // removed or filtered
     CloseableIterable<TrackedFile> files =
         CloseableIterable.transform(open(), this::applyInheritance);
 
     if (!includeAll) {
-      files = CloseableIterable.filter(files, file -> 
file.tracking().isLive());
+      if (dv != null) {
+        files =
+            CloseableIterable.filter(files, file -> file.tracking().isLive() 
&& !isDeleted(file));
+      } else {
+        files = CloseableIterable.filter(files, file -> 
file.tracking().isLive());
+      }
+    } else if (dv != null) {
+      files = CloseableIterable.transform(files, this::applyMDVDeletes);

Review Comment:
   Thanks for the context in the PR description.
   
   I am wondering if we should just fail here (instead of returning partially 
correct result)? or remove includeAll for now (like you mentioned) until we 
have the proper support for change detection read?



##########
core/src/main/java/org/apache/iceberg/V4ManifestReader.java:
##########
@@ -156,6 +167,20 @@ private TrackedFile applyInheritance(TrackedFile file) {
     return file;
   }
 
+  private boolean isDeleted(TrackedFile file) {

Review Comment:
   nit: my initial reaction of this method name is the status `deleted`. should 
we change it to `isPositionDeleted`?



##########
core/src/main/java/org/apache/iceberg/TrackingStruct.java:
##########
@@ -179,6 +180,32 @@ boolean assignFirstRowId(Long nextRowId) {
     return false;
   }
 
+  /**
+   * Update this tracking to have DELETED status.
+   *
+   * @param snapshotId snapshot ID when the record was deleted, or null if it 
is unknown
+   * @return this converted to a DELETED Tracking
+   */
+  TrackingStruct convertToDeleted(Long snapshotId) {
+    Preconditions.checkState(status.isLive(), "Cannot convert %s to DELETED", 
status);
+    this.status = EntryStatus.DELETED;
+    this.snapshotId = snapshotId;
+    return this;
+  }
+
+  /**
+   * Update this tracking to have REPLACED status.
+   *
+   * @param snapshotId snapshot ID when the record was replaced, or null if it 
is unknown
+   * @return this converted to a REPLACED Tracking
+   */
+  TrackingStruct convertToReplaced(Long snapshotId) {

Review Comment:
   this method is not used in this PR.



##########
core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java:
##########
@@ -370,146 +467,178 @@ public void 
inheritanceUncommittedOnlyInheritsSnapshotId(FileFormat format) thro
             .metricsConfig(METRICS_CONFIG);
     List<TrackedFile> actual = read(builder);
 
+    List<TrackedFile> expected =
+        List.of(
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.ADDED, 1234567L, null, null, 
null, null, null, null),
+                "s3://bucket/table/file-b.parquet"),
+            unpartitionedDataFile( // inherits only snapshot ID
+                new TrackingStruct(EntryStatus.ADDED, 34L, null, null, null, 
null, null, null),
+                "s3://bucket/table/file-a.parquet"));
+
     assertThat(actual)
-        .extracting(file -> file.tracking().snapshotId())
-        .containsExactly(1234567L, SNAPSHOT_ID);
-    assertThat(actual)
-        .as("Uncommitted files do not inherit data sequence number")
-        .extracting(file -> file.tracking().dataSequenceNumber())
-        .containsOnlyNulls();
-    assertThat(actual)
-        .as("Uncommitted files do not inherit file sequence number")
-        .extracting(file -> file.tracking().fileSequenceNumber())
-        .containsOnlyNulls();
+        .usingComparatorForType(FILE_AND_TRACKING_COMPARATOR, 
TrackedFile.class)
+        .isEqualTo(expected);
   }
 
   @ParameterizedTest
   @FieldSource("MANIFEST_FORMATS")
   public void inheritanceSnapshotId(FileFormat format) throws IOException {
-    Tracking trackingWithSnapshotId =
-        new TrackingStruct(EntryStatus.ADDED, 1234567L, null, null, null, 
null, null, null);
     TrackedFile withSnapshotId =
-        unpartitionedDataFile(trackingWithSnapshotId, 
"s3://bucket/table/file-b.parquet");
-    Tracking trackingWithoutSnapshotId =
-        new TrackingStruct(EntryStatus.ADDED, null, null, null, null, null, 
null, null);
+        unpartitionedDataFile(
+            new TrackingStruct(EntryStatus.ADDED, 1234567L, 5L, 5L, null, 
5_000L, null, null),
+            "s3://bucket/table/file-b.parquet");
     TrackedFile withoutSnapshotId =
-        unpartitionedDataFile(trackingWithoutSnapshotId, 
"s3://bucket/table/file-a.parquet");
+        unpartitionedDataFile(
+            new TrackingStruct(EntryStatus.ADDED, null, 5L, 5L, null, 5_100L, 
null, null),
+            "s3://bucket/table/file-a.parquet");
 
     ManifestFile manifest =
         writeManifest(
             format, UNPARTITIONED_TYPE, ImmutableList.of(withSnapshotId, 
withoutSnapshotId));
 
+    when(manifest.firstRowId()).thenReturn(10_000L);
+    when(manifest.snapshotId()).thenReturn(34L);
+
     V4ManifestReader.Builder builder =
         V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, 
ID_PARTITIONING_SPECS)
             .metricsConfig(METRICS_CONFIG);
     List<TrackedFile> actual = read(builder);
 
+    List<TrackedFile> expected =
+        List.of(
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.ADDED, 1234567L, 5L, 5L, null, 
5_000L, null, null),
+                "s3://bucket/table/file-b.parquet"),
+            unpartitionedDataFile( // inherits only snapshot ID
+                new TrackingStruct(EntryStatus.ADDED, 34L, 5L, 5L, null, 
5_100L, null, null),
+                "s3://bucket/table/file-a.parquet"));
+
     assertThat(actual)
-        .extracting(file -> file.tracking().snapshotId())
-        .containsExactly(1234567L, SNAPSHOT_ID);
+        .usingComparatorForType(FILE_AND_TRACKING_COMPARATOR, 
TrackedFile.class)
+        .isEqualTo(expected);
   }
 
   @ParameterizedTest
   @FieldSource("MANIFEST_FORMATS")
   public void inheritanceDVSnapshotIdNotInherited(FileFormat format) throws 
IOException {
-    Tracking trackingWithDVSnapshotId =
-        new TrackingStruct(EntryStatus.ADDED, SNAPSHOT_ID, null, null, 
1234567L, null, null, null);
     TrackedFile withDVSnapshotId =
-        unpartitionedDataFile(trackingWithDVSnapshotId, 
"s3://bucket/table/file-b.parquet");
-    Tracking trackingWithoutDVSnapshotId =
-        new TrackingStruct(EntryStatus.ADDED, SNAPSHOT_ID, null, null, null, 
null, null, null);
+        unpartitionedDataFile(
+            new TrackingStruct(
+                EntryStatus.ADDED, SNAPSHOT_ID, 5L, 5L, 1234567L, 5_000L, 
null, null),
+            "s3://bucket/table/file-b.parquet");
     TrackedFile withoutDVSnapshotId =
-        unpartitionedDataFile(trackingWithoutDVSnapshotId, 
"s3://bucket/table/file-a.parquet");
+        unpartitionedDataFile(
+            new TrackingStruct(EntryStatus.ADDED, SNAPSHOT_ID, 5L, 5L, null, 
5_100L, null, null),
+            "s3://bucket/table/file-a.parquet");
 
     ManifestFile manifest =
         writeManifest(
             format, UNPARTITIONED_TYPE, ImmutableList.of(withDVSnapshotId, 
withoutDVSnapshotId));
 
+    when(manifest.firstRowId()).thenReturn(10_000L);
+    when(manifest.snapshotId()).thenReturn(34L);
+
     V4ManifestReader.Builder builder =
         V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, 
ID_PARTITIONING_SPECS)
             .metricsConfig(METRICS_CONFIG);
     List<TrackedFile> actual = read(builder);
 
+    List<TrackedFile> expected =
+        List.of(
+            unpartitionedDataFile(
+                new TrackingStruct(
+                    EntryStatus.ADDED, SNAPSHOT_ID, 5L, 5L, 1234567L, 5_000L, 
null, null),
+                "s3://bucket/table/file-b.parquet"),
+            unpartitionedDataFile( // inherits only snapshot ID

Review Comment:
   nit: the comment probably should be `// dv_snapshot_id stays null (not 
inherited)`



##########
core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java:
##########
@@ -318,49 +321,143 @@ public void readManifestFile(FileFormat format) throws 
IOException {
   @FieldSource("MANIFEST_FORMATS")
   public void statusFilter(FileFormat format) throws IOException {
     List<TrackedFile> files =
-        ImmutableList.of(
-            unpartitionedFileWithStatus(EntryStatus.ADDED, 
"s3://bucket/added.parquet"),
-            unpartitionedFileWithStatus(EntryStatus.MODIFIED, 
"s3://bucket/modified.parquet"),
-            unpartitionedFileWithStatus(EntryStatus.DELETED, 
"s3://bucket/deleted.parquet"),
-            unpartitionedFileWithStatus(EntryStatus.EXISTING, 
"s3://bucket/existing.parquet"),
-            unpartitionedFileWithStatus(EntryStatus.REPLACED, 
"s3://bucket/replaced.parquet"));
+        List.of(
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.ADDED, 42L, null, null, null, 
null, null, null),
+                "s3://bucket/added.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.MODIFIED, 40L, 5L, 5L, 42L, 
5_000L, null, null),
+                "s3://bucket/modified.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.DELETED, 42L, 2L, 2L, null, 
1_000L, null, null),
+                "s3://bucket/deleted.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.EXISTING, 38L, 4L, 4L, null, 
3_000L, null, null),
+                "s3://bucket/existing.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.REPLACED, 42L, 2L, 2L, 40L, 
2_000L, null, null),
+                "s3://bucket/replaced.parquet"));
 
     ManifestFile manifest = writeManifest(format, UNPARTITIONED_TYPE, files);
+    when(manifest.firstRowId()).thenReturn(10_000L);
 
-    List<TrackedFile> liveFiles =
-        read(
-            V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, 
UNPARTITIONED_SPECS)
-                .metricsConfig(METRICS_CONFIG));
+    List<TrackedFile> expectedFiles =
+        List.of(
+            unpartitionedDataFile( // inherits seq number and assigned first 
row ID
+                new TrackingStruct(
+                    EntryStatus.ADDED, 42L, MANIFEST_SEQ, MANIFEST_SEQ, null, 
10_000L, null, null),
+                "s3://bucket/added.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.MODIFIED, 40L, 5L, 5L, 42L, 
5_000L, null, null),
+                "s3://bucket/modified.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.DELETED, 42L, 2L, 2L, null, 
1_000L, null, null),
+                "s3://bucket/deleted.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.EXISTING, 38L, 4L, 4L, null, 
3_000L, null, null),
+                "s3://bucket/existing.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.REPLACED, 42L, 2L, 2L, 40L, 
2_000L, null, null),
+                "s3://bucket/replaced.parquet"));
+
+    V4ManifestReader.Builder builder =
+        V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, 
UNPARTITIONED_SPECS)
+            .metricsConfig(METRICS_CONFIG);
+
+    List<TrackedFile> liveFiles = read(builder);
     assertThat(liveFiles)
-        .usingComparatorForType(FILE_COMPARATOR, TrackedFile.class)
-        .containsExactly(files.get(0), files.get(1), files.get(3));
+        .usingComparatorForType(FILE_AND_TRACKING_COMPARATOR, 
TrackedFile.class)
+        .containsExactly(expectedFiles.get(0), expectedFiles.get(1), 
expectedFiles.get(3));
 
-    List<TrackedFile> allFiles =
-        read(
-            V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, 
UNPARTITIONED_SPECS)
-                .metricsConfig(METRICS_CONFIG)
-                .includeAll());
+    List<TrackedFile> allFiles = read(builder.includeAll());
     assertThat(allFiles)
-        .usingComparatorForType(FILE_COMPARATOR, TrackedFile.class)
-        .containsExactlyElementsOf(files);
+        .usingComparatorForType(FILE_AND_TRACKING_COMPARATOR, 
TrackedFile.class)
+        .isEqualTo(expectedFiles);
+  }
+
+  @ParameterizedTest
+  @FieldSource("MANIFEST_FORMATS")
+  public void mdvFilter(FileFormat format) throws IOException {
+    List<TrackedFile> files =
+        List.of(
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.ADDED, 42L, null, null, null, 
null, null, null),
+                "s3://bucket/added.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.MODIFIED, 40L, 5L, 5L, 42L, 
5_000L, null, null),
+                "s3://bucket/modified.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.DELETED, 42L, 2L, 2L, null, 
1_000L, null, null),
+                "s3://bucket/deleted.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.EXISTING, 38L, 4L, 4L, null, 
3_000L, null, null),
+                "s3://bucket/existing.parquet"),
+            unpartitionedDataFile(
+                new TrackingStruct(EntryStatus.REPLACED, 42L, 2L, 2L, 40L, 
2_000L, null, null),
+                "s3://bucket/replaced.parquet"));
+
+    ManifestFile manifest = writeManifest(format, UNPARTITIONED_TYPE, files);
+
+    // this bitmap deletes a MODIFIED entry and an already DELETED entry
+    when(manifest.manifestDeletionVector())
+        .thenReturn(MumblingTestUtil.bitmap(MumblingTestUtil.sparse(1, 2)));
+    when(manifest.firstRowId()).thenReturn(10_000L);
+
+    List<TrackedFile> expectedFiles =
+        List.of(
+            unpartitionedDataFile( // inherits seq number and assigned first 
row ID
+                new TrackingStruct(
+                    EntryStatus.ADDED, 42L, MANIFEST_SEQ, MANIFEST_SEQ, null, 
10_000L, null, null),
+                "s3://bucket/added.parquet"),
+            unpartitionedDataFile( // status changed to DELETED, delete 
snapshot ID is unknown

Review Comment:
   why is the modified entry converted to deleted?



##########
core/src/test/java/org/apache/iceberg/V4TestComparators.java:
##########
@@ -94,9 +94,7 @@ private static <T extends Comparable<T>> Comparator<T> 
natural() {
               .thenComparing(Tracking::dvSnapshotId, natural())
               .thenComparing(Tracking::firstRowId, natural())
               .thenComparing(Tracking::deletedPositions, BYTES)
-              .thenComparing(Tracking::replacedPositions, BYTES)
-              .thenComparing(Tracking::manifestLocation, natural())
-              .thenComparingLong(Tracking::manifestPos));
+              .thenComparing(Tracking::replacedPositions, BYTES));

Review Comment:
   should we add the context as inline code comment for future reference on why 
the two fields are not included in the comparator construction?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to