JingsongLi commented on code in PR #10105:
URL: https://github.com/apache/paimon/pull/10105#discussion_r4092281796
##########
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:
[P1] Avoid using the overwrite timestamp as the snapshot-branch boundary.
`earlierOrEqualTimeMills` can select a snapshot-branch commit made *after*
the delta overwrite. The branches have independent histories, and their
wall-clock timestamps have millisecond precision and can differ across hosts. I
reproduced this with a focused `ChainTableFileStoreTableTest`: publish the
delta overwrite while its snapshot-branch cleanup fails, append a new row to
the snapshot branch, then force only that later snapshot's recorded
`timeMillis` to equal the overwrite's. On replay, this line selects the later
snapshot as `asOf`, and the cleanup deletes the new row along with the
superseded row; the read falls back to the delta row. The existing
`testChainOverwriteReplayClearsOnlyWhatTheOverwriteSuperseded` passes with
separated timestamps, while this equal-timestamp case fails.
Please use an exact snapshot-branch position associated with the overwrite,
or fail closed when that position cannot be established. A wall-clock
comparison across branches cannot safely decide which files the overwrite
superseded.
--
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]