zhuxiangyi opened a new pull request, #10204:
URL: https://github.com/apache/paimon/pull/10204
### Purpose
Fix a data-loss path in Flink's `rollback_to_as_latest` when another writer
commits between the procedure's initial snapshot lookup and the core rollback,
followed by a post-commit failure.
When invoked with `snapshot_id`, the procedure creates a protection tag so
that subsequent snapshot expiration does not delete files restored from the
target snapshot. Its exception handler currently guesses the rollback snapshot
ID using the **initial** `latestSnapshot.id() + 1`, then checks that snapshot's
`commitUser` to decide whether to delete the tag. However, the core rollback
reads the latest snapshot again, and an exception does not necessarily mean
that the snapshot was not committed.
For example:
1. The procedure reads latest snapshot **3** and creates a protection tag
for target snapshot **1**.
2. Another writer commits snapshot **4** before the core rollback reads
latest.
3. The rollback successfully publishes snapshot **5**, restoring snapshot
1's data.
4. A post-commit callback throws.
5. Cleanup checks snapshot **4**, finds another writer's `commitUser`, and
incorrectly deletes the protection tag.
6. Snapshot 5 is initially readable, but subsequent expiration of the older
snapshots can delete its restored data files, causing `FileNotFoundException`.
The interleaved writer finishes before the core rollback reads latest, so a
snapshot commit conflict is not required for this sequence.
#### Fix
Determine whether cleanup is safe from the actual call boundary and return
value, rather than a predicted snapshot ID:
- Before entering core rollback, failures may clean up the newly created tag.
- If core rollback explicitly returns `false`, cleanup is allowed.
- If core rollback throws, retain the protection tag: the snapshot might
already be committed.
- Successful rollback continues to retain the tag. User-supplied tags are
not deleted.
Remove the stale snapshot-ID/commit-user lookup helpers. No core or Spark
code is changed.
This deliberately favors data safety over reclaiming a tag after an
ambiguous failure. An exception inside core rollback that actually occurs
before publication may also leave a protection tag, because the procedure
cannot safely distinguish that case from a committed rollback followed by a
failure.
### Tests
Add regression coverage in `RollbackToAsLatestProcedureITCase`:
- Exercise the real Flink SQL `CALL` with a post-commit callback failure,
both with and without an interleaved real commit. Assert that the protection
tag remains and the restored data is readable before and after snapshot
expiration.
- Verify cleanup when commit construction fails and when rollback explicitly
returns `false`.
The interleaving and callback exception are injected deterministically; this
is not a probabilistic concurrency test or a live metastore outage. The SQL
procedure, successful commits, tag lifecycle, expiration, and reads use the
real implementations. Known-not-committed cleanup cases use mocked commit
outcomes.
Verified on a branch based on `master` at `e42cd1302`:
- With the original production implementation and the new tests: **4 cases,
1 failure**, because the interleaved callback-failure case loses its protection
tag.
- With the fix, the entire procedure test class: **10 cases, 0 failures, 0
errors**.
- A separate local Flink SQL reproduction that previously reached
`FileNotFoundException` also remains readable after expiration with this fix.
```bash
mvn -pl paimon-flink/paimon-flink-common -Pflink1 \
-DwildcardSuites=none \
-Dtest=RollbackToAsLatestProcedureITCase test
git diff --check
```
The final Maven run did not use `fast-build`; Checkstyle, Spotless, and
Enforcer checks passed. Tested with Flink 1.20.1 on JDK 17; Flink 2 and the
full reactor were not run.
--
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]