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]

Reply via email to