JingsongLi commented on PR #10198:
URL: https://github.com/apache/paimon/pull/10198#issuecomment-5950888416

   [P2] Prefer the earlier tagged snapshot when watermarks tie 
(`CreateTagFromWatermarkProcedure.java:98–101`)
   
   The new procedure promises the first qualifying snapshot, including tagged 
snapshots, but the strict `tagSnapshot.watermark() < snapshot.watermark()` 
comparison ignores an older tagged snapshot with the same watermark as the 
retained candidate. Expiration can therefore change the selected dataset even 
though the earlier snapshot is still available through a tag.
   
   I reproduced this with real stream commits, snapshot expiration, Spark SQL 
CALL, and tag time-travel reads:
   
   - Snapshot 1: watermark 1000, row `(1, 'row1')`; preserve it as tag 
`historical`.
   - Snapshot 2: watermark 1000, with the additional row `(2, 'row2')`.
   - Expire snapshot 1 while keeping snapshot 2. The historical tag still reads 
only row 1.
   - `CALL ...create_tag_from_watermark(table => 'test.T', tag => 'boundary', 
watermark => 500)` creates a tag for snapshot 2. Reading `T VERSION AS OF 
'boundary'` returns rows 1 and 2, instead of the first qualifying snapshot's 
row 1.
   
   `SnapshotManager` correctly returns the first *retained* candidate, snapshot 
2. This procedure's new merge with tagged snapshots should select snapshot 1, 
but equal watermarks fail the strict comparison. Repeated watermarks are 
normal: commits retain the maximum/inherit the previous watermark, so a stalled 
Flink watermark or a Spark append after a Flink commit can produce this state.
   
   When candidate watermarks are equal, choose the smaller snapshot ID. Simply 
changing `<` to `<=` would instead let a newer tag replace an earlier retained 
snapshot, so cover that direction as well. The Flink counterpart has the same 
condition, but this new Spark API should meet its documented first-snapshot 
contract rather than propagate that defect.
   
   Validation: the wrong tagged dataset reproduces on both Spark 3.5 and Spark 
4.1. On Spark 3.5, the same request selects snapshot 1 before expiration and 
snapshot 2 afterward; the control with only a newer equal-watermark tag 
correctly keeps earlier retained snapshot 1. The PR's original actual SQL suite 
passes on both versions with normal Maven checks, and the 1-day tag-retention 
argument is verified in the reproduction.
   


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