Repository navigation
[spark] Support create_tag_from_watermark procedure - #10420
Stephen0421 wants to merge 1 commit into
Conversation
eee40f2 to
6118039
Compare
|
[P2] Choose the earlier snapshot when a retained tag has the same watermark The tag candidate is compared to the live candidate only with I reproduced this through real Spark 3 SQL:
The new tag is created on snapshot 2, rather than the still-retained tagged snapshot 1. These snapshots contain different data. Repeated watermarks are normal: Please break equal-watermark ties by snapshot ID, preferring the earlier snapshot. Simply changing Validation: both original procedure tests passed with standard JDK 8/Spark 3 Maven checks. Adding the expired-tag/equal-watermark SQL case produced two passes and one failure ( |
6118039 to
4606c1c
Compare
Fixed. When watermarks are equal, the procedure now keeps the earlier snapshot by comparing snapshot IDs, instead of changing Two tests cover both directions:
This change is limited to the new Spark procedure. The Flink implementation still uses the original comparison. |
Purpose
Add the Spark procedure
sys.create_tag_from_watermark, matching the existing Flink procedure.The procedure finds the earliest retained snapshot whose watermark is greater than or equal to the given value, including snapshots that have already expired but are still kept by a tag, then creates a tag on that snapshot. Snapshots with a null watermark are skipped. If none matches, it throws
SnapshotNotExistException.This only reads watermarks already stored on snapshots (for example from Flink or Spark Structured Streaming writes). It does not generate a watermark.
Output columns:
tagName,snapshot,commit_time,watermark.Test
CreateTagFromWatermarkProcedureTest: tag the first snapshot at or after the watermark, and fall back to an expired snapshot that is still kept by a tagSnapshotNotExistException