Skip to content

[spark] Add rollback_to_as_latest procedure - #10199

Merged
JingsongLi merged 4 commits into
apache:masterfrom
jackylee-ch:spark-rollback-to-as-latest
Oct 3, 2026
Merged

JingsongLi merged 4 commits into
apache:masterfrom
jackylee-ch:spark-rollback-to-as-latest

Conversation

@jackylee-ch

Copy link
Copy Markdown
Contributor

Purpose

Spark's rollback drops every snapshot and tag created after the target, so a
user who only wants to restore an earlier state loses the newer history. Flink
already exposes rollback_to_as_latest, which rolls the target forward as a
new latest snapshot non-destructively (core rollbackToAsLatest, #8450), but
Spark has no equivalent.

Change

Add RollbackToAsLatestProcedure, registered as rollback_to_as_latest,
taking exactly one of tag / snapshot_id. It mirrors the Flink procedure:
resolve the target snapshot (creating a protection tag for the snapshot_id
path), commit.rollbackToAsLatest(tag), and clean that tag up if the commit
did not land. Returns previous / rolled_back / current snapshot ids.

Tests

RollbackToAsLatestProcedureTest writes three snapshots, rolls back to
snapshot 1 as latest and asserts snapshots 2 and 3 survive (non-destructive)
and the data is snapshot 1's, then rolls forward to snapshot 3.

Written with Claude Code; verification is mine.

Spark's rollback drops every snapshot and tag created after the target, so a
user who only wants to restore an earlier state loses the newer history. Flink
already exposes rollback_to_as_latest, which rolls the target forward as a new
latest snapshot non-destructively (core rollbackToAsLatest, apache#8450), but Spark
has no equivalent.

Add RollbackToAsLatestProcedure, registered as rollback_to_as_latest, taking
exactly one of tag / snapshot_id. It mirrors the Flink procedure: resolve the
target snapshot (creating a protection tag for the snapshot_id path),
commit.rollbackToAsLatest(tag), and clean that tag up if the commit did not
land. Returns previous / rolled_back / current snapshot ids.

Tests: RollbackToAsLatestProcedureTest writes three snapshots, rolls back to
snapshot 1 as latest and asserts snapshots 2 and 3 survive (non-destructive)
and the data is snapshot 1's, then rolls forward to snapshot 3.

Written with Claude Code; verification is mine.

@zhuxiangyi zhuxiangyi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the Spark procedure. I found one data-safety issue in the failure cleanup path and reproduced it through Spark SQL; details are inline. Could we address this before merging?

store,
tagManager,
createdRollbackTag,
latestSnapshot.id() + 1,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we avoid using the previously read latestSnapshot.id() + 1 to decide whether the rollback committed? FileStoreCommitImpl.rollbackToAsLatest() reads the latest snapshot again, so another writer can commit between these two reads without causing a commit conflict.

For example:

  1. This procedure reads latest = 3.
  2. Another writer commits snapshot 4.
  3. The rollback successfully commits snapshot 5.
  4. A post-commit callback throws.
  5. Cleanup checks snapshot 4, sees a different commitUser, and deletes the protection tag even though snapshot 5 committed successfully.

I reproduced this through the Spark SQL procedure using a deterministically interleaved real commit and an injected post-commit callback failure. Snapshot 5 was initially readable, but the protection tag was gone; expiring the older snapshots then deleted the restored data file, and SELECT failed with FileNotFoundException. The control case with the same callback failure but no interleaved commit retained the tag and remained readable after expiration.

Please establish the actual commit outcome rather than infer it from this stale snapshot ID; when the outcome is uncertain, retain the protection tag. A regression test should cover the interleaved commit + post-commit failure + subsequent snapshot expiration case.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, thank you — confirmed. Cleanup no longer infers the outcome from latestSnapshot.id() + 1; it now scans every snapshot committed after the initial read for one authored by this rollback's unique commit user, and deletes the protection tag only when it positively confirms the rollback did not commit — if it committed or the outcome is uncertain, the tag is retained (so your snapshot-5 case keeps it). The mirrored Flink procedure had the same flaw and is fixed identically. Added a regression test with distinct commit users on snapshots 4 and 5.

… after an interleaved commit

Data-safety fix in the failure-path cleanup of rollback_to_as_latest.

The procedure reads the latest snapshot, creates a protection tag for the
rollback target, then commits the rollback. On any exception (e.g. a
post-commit callback throwing) it decided whether the rollback had committed
by probing snapshot `latestSnapshot.id() + 1`. But
FileStoreCommitImpl.rollbackToAsLatest re-reads the latest snapshot, so a
writer committing between the two reads pushes the rollback to a higher id:

1. procedure reads latest = 3
2. another writer commits snapshot 4
3. the rollback commits snapshot 5
4. a post-commit callback throws
5. cleanup probes snapshot 4, sees a different commitUser, and deletes the
   protection tag though snapshot 5 committed -- expiring older snapshots
   then drops the restored files (FileNotFoundException on read).

Fix: establish the real outcome instead of inferring it from the stale id.
Cleanup now scans every snapshot committed after the initial read for one
authored by this rollback's unique commit user, and only deletes the tag
when it positively confirms the rollback did not commit; if it committed --
or the outcome is uncertain -- the tag is retained. The mirrored Flink
procedure had the identical flaw and is fixed the same way.

Tests: a regression test drives two rollbacks so snapshots 4 and 5 carry
distinct commit users, then asserts the cleanup check finds the rollback at
5 rather than probing the stale id 4 (the old point-probe returned false).

Written with Claude Code; verification is mine.
@jackylee-ch
jackylee-ch force-pushed the spark-rollback-to-as-latest branch from e1bf465 to 123273e Compare September 26, 2026 15:00
…op the Flink change

The failure cleanup previously inferred the rollback outcome from a snapshot
id, which is unsafe when another writer commits between the procedure's
snapshot read and the core rollback. Track whether cleanup is safe from the
call boundary instead: only delete the created protection tag when we never
entered the commit or the commit explicitly returned false; if the commit
threw after publishing (e.g. a post-commit callback), keep the tag so
expiration cannot drop the restored files. This mirrors the Flink-side
approach in apache#10204, so the Flink change is dropped here and left to that PR.
@jackylee-ch

Copy link
Copy Markdown
Contributor Author

Thanks @zhuxiangyi. Reworked to mirror your #10204 approach instead of probing a snapshot id: cleanup deletes the protection tag only when the commit was never entered or returned false; if it threw after publishing (e.g. a post-commit callback), the tag is kept. I also dropped the Flink change so #10204 owns it — this PR is Spark-only now.

Verified in paimon-spark-ut: new tests run the real CALL with a callback failure, with and without an interleaved commit, and assert the tag survives expiration; the old latestSnapshot.id() + 1 cleanup fails the interleaved case. Added Spark docs; spotless clean.

@JingsongLi

Copy link
Copy Markdown
Contributor

[P2] Refresh Spark's cached table after an ambiguously successful rollback (RollbackToAsLatestProcedure.java:161–165)

The new failure path preserves the protection tag correctly, but it rethrows before BaseProcedure.modifyPaimonTable reaches refreshSparkCache. If a post-commit callback throws, the rollback snapshot is already durable while CACHE TABLE continues serving the pre-rollback data.

I reproduced this through real Spark SQL using the PR's FailingRollbackCallback:

  1. Write snapshot 1 with (1, 'original'), then overwrite with snapshot 2 containing (2, 'replacement').
  2. Run CACHE TABLE T and read the replacement row.
  3. Set failRollbackCommit = true, then call rollback_to_as_latest(..., snapshot_id => 1). The call reports the injected post-commit callback error, and snapshot 3 has nevertheless been published.
  4. SELECT * FROM T still returns (2, 'replacement'). After UNCACHE TABLE T, the identical query returns (1, 'original') from the persisted rollback.

The shared helper's success-only refresh predates this PR, but this new procedure introduces the public rollback failure path that reaches it; the issue is in this wrapper's handling of the already-published state. The repository's MSCK partial-failure tests similarly require cached reads to reflect durable changes even when the command reports an error.

Please refresh/invalidate cached plans after an entered rollback with an uncertain outcome, including this exception path, and preserve the original rollback exception if refreshing also fails. Add a CACHE TABLE regression alongside the existing post-commit callback tests. The successful rollback path already refreshes correctly.

Verified on both Spark 3.5 and Spark 4.1: the cached-vs-persisted mismatch reproduces on each. All three original rollback tests pass on both with normal Maven checks; additional real CALL checks for exclusive/missing arguments, tag rollback after snapshot expiration, snapshot-id fallback through a retained tag, and successful cache refresh also pass.

@jackylee-ch

Copy link
Copy Markdown
Contributor Author

Fixed in ab2a2d2. The procedure now refreshes Spark's cached plans even when the rollback fails after entering the commit — rollback_to_as_latest can publish the rollback snapshot and then throw (e.g. a post-commit callback), so the state is already durable and the shared success-only refresh left CACHE TABLE serving pre-rollback data. A small wrapper runs refreshSparkCache in both the success and failure paths; the original failure is preserved and a refresh error is only added as suppressed.

Verified locally in paimon-spark-ut (spark3, 4 passed): the new test caches the table, injects a post-commit callback failure during the rollback, and asserts SELECT returns the rolled-back row rather than the stale cached one. Reverting to the plain wrapper fails it with [1,original] vs [2,replacement]. spotless and git diff --check clean.

@JingsongLi

Copy link
Copy Markdown
Contributor

+1

@JingsongLi
JingsongLi merged commit 283f709 into apache:master Oct 3, 2026
11 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants