zhuxiangyi commented on code in PR #10105:
URL: https://github.com/apache/paimon/pull/10105#discussion_r4094559620


##########
paimon-core/src/main/java/org/apache/paimon/metastore/ChainTableOverwriteCommitCallback.java:
##########
@@ -52,52 +73,305 @@
  */
 public class ChainTableOverwriteCommitCallback implements CommitCallback {
 
+    private static final Logger LOG =
+            LoggerFactory.getLogger(ChainTableOverwriteCommitCallback.class);
+
     private transient FileStoreTable table;
     private transient CoreOptions coreOptions;
+    private final String commitUser;
 
-    public ChainTableOverwriteCommitCallback(FileStoreTable table) {
+    public ChainTableOverwriteCommitCallback(FileStoreTable table, String 
commitUser) {
         this.table = table;
         this.coreOptions = table.coreOptions();
+        this.commitUser = commitUser;
     }
 
     @Override
     public void call(Context context) {
+        if (!ChainTableUtils.isScanFallbackDeltaBranch(coreOptions)) {
+            return;
+        }
+        if (context.snapshot.commitKind() != CommitKind.OVERWRITE) {
+            return;
+        }
+        truncateSnapshotPartitions(context.deltaFiles);
+    }
 
+    /**
+     * The commit of this committable was published by an earlier attempt 
whose callback may not
+     * have completed, for example because the snapshot branch was unreachable 
right after the delta
+     * snapshot was written. Resolve that snapshot and redo the cleanup, which 
is idempotent. The
+     * partitions are taken from the manifest changes of the snapshot rather 
than from the
+     * committable, since an overwrite also clears partitions it wrote no new 
file to.
+     */
+    @Override
+    public void retry(ManifestCommittable committable) {
         if (!ChainTableUtils.isScanFallbackDeltaBranch(coreOptions)) {
             return;
         }
+        List<Snapshot> snapshots =
+                table.snapshotManager()
+                        .findSnapshotsForIdentifiers(
+                                commitUser, 
Collections.singletonList(committable.identifier()));
+        if (snapshots.isEmpty()) {
+            LOG.warn(
+                    "No snapshot of commit user {} with identifier {} in table 
{}, "
+                            + "cannot redo the snapshot branch cleanup of its 
overwrite.",
+                    commitUser,
+                    committable.identifier(),
+                    table.name());
+            return;
+        }
+        for (Snapshot snapshot : snapshots) {
+            if (snapshot.commitKind() != CommitKind.OVERWRITE) {
+                continue;
+            }
+            List<BinaryRow> overwritePartitions =
+                    overwritePartitions(
+                            table.store()
+                                    .newScan()
+                                    .withKind(ScanMode.DELTA)
+                                    .withSnapshot(snapshot.id())
+                                    .plan()
+                                    .files());
+            clearSnapshotFilesAsOf(overwritePartitions, snapshot);
+        }
+    }
 
-        if (context.snapshot.commitKind() != CommitKind.OVERWRITE) {
+    /**
+     * Clear what the overwrite superseded in the given partitions of the 
snapshot branch, and
+     * nothing that landed there since. That is what {@link #call} cleared at 
the time; a replay
+     * that repeats it after later data arrived must not take that data with 
it, whether the
+     * original cleanup had completed or not.
+     *
+     * <p>What the overwrite superseded are the files the snapshot branch held 
in those partitions
+     * when the overwrite was published, and whatever compactions of the 
snapshot branch have since
+     * rewritten from them alone. If a compaction has merged them with data 
written after the
+     * overwrite, another commit such as a rescale has rewritten them, or the 
snapshot branch no
+     * longer retains the history to tell, the cleanup cannot be done exactly, 
and the retry fails
+     * rather than report a cleanup it did not do.
+     */
+    private void clearSnapshotFilesAsOf(List<BinaryRow> partitions, Snapshot 
overwrite) {
+        if (partitions.isEmpty()) {
+            return;
+        }
+        FileStoreTable snapshotTable = snapshotTable();
+        SnapshotManager snapshotManager = snapshotTable.snapshotManager();
+        Snapshot latest = snapshotManager.latestSnapshot();
+        if (latest == null) {
             return;
         }
+        Snapshot asOf = 
snapshotManager.earlierOrEqualTimeMills(overwrite.timeMillis());

Review Comment:
   Fixed in 27290b3: the overwrite now records where the snapshot branch stands 
in a property of its snapshot, through a new `CommitCallback#beforeOverwrite` 
hook called before it is published, and `retry` takes that snapshot as the 
boundary; the time lookup is gone. It fails closed when no position was 
recorded or the recorded snapshot has expired, and an unreadable snapshot 
branch fails the overwrite before publishing. Your equal-timestamp case is 
`testChainOverwriteReplayWithSnapshotBranchCommitAtTheSameTime`. Details in 
https://github.com/apache/paimon/pull/10105#issuecomment-5815663818.
   



-- 
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]

Reply via email to