Skip to content

[spark] Support create_tag_from_watermark procedure - #10420

Open
Stephen0421 wants to merge 1 commit into
apache:masterfrom
Stephen0421:spark-create-tag-from-watermark
Open

Stephen0421 wants to merge 1 commit into
apache:masterfrom
Stephen0421:spark-create-tag-from-watermark

Conversation

@Stephen0421

Copy link
Copy Markdown
Contributor

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.

CALL sys.create_tag_from_watermark(
  `table` => 'default.T',
  `tag` => 'my_tag',
  `watermark` => 1724404318750,
  time_retained => '1 d'
);

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 tag
  • A watermark later than every snapshot throws SnapshotNotExistException

@Stephen0421
Stephen0421 force-pushed the spark-create-tag-from-watermark branch from eee40f2 to 6118039 Compare October 8, 2026 04:05
@JingsongLi

Copy link
Copy Markdown
Contributor

[P2] Choose the earlier snapshot when a retained tag has the same watermark

The tag candidate is compared to the live candidate only with tagSnapshot.watermark() < snapshot.watermark() in CreateTagFromWatermarkProcedure. Equal watermarks therefore keep the later live snapshot, although the new API promises the first retained snapshot, including tagged snapshots.

I reproduced this through real Spark 3 SQL:

  1. Write snapshot 1 with watermark 1000 and create keep1 on snapshot 1.
  2. Write additional data in snapshot 2 with the same watermark 1000.
  3. Expire snapshot 1 while retaining keep1.
  4. Call sys.create_tag_from_watermark(..., watermark => 1000).

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: FileStoreCommitImpl inherits the previous value for null watermarks and takes the maximum for new values.

Please break equal-watermark ties by snapshot ID, preferring the earlier snapshot. Simply changing < to <= would choose a later tagged snapshot in the reverse arrangement, so please cover both directions. The existing Flink implementation shares this comparison; this finding concerns the new Spark API and its stated contract, rather than a Flink regression.

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 (expected snapshot=1, actual=2); independent review confirmed the reachable watermark/tags ordering semantics.

@Stephen0421
Stephen0421 force-pushed the spark-create-tag-from-watermark branch from 6118039 to 4606c1c Compare October 10, 2026 06:07
@Stephen0421

Copy link
Copy Markdown
Contributor Author

[P2] Choose the earlier snapshot when a retained tag has the same watermark

The tag candidate is compared to the live candidate only with tagSnapshot.watermark() < snapshot.watermark() in CreateTagFromWatermarkProcedure. Equal watermarks therefore keep the later live snapshot, although the new API promises the first retained snapshot, including tagged snapshots.

I reproduced this through real Spark 3 SQL:

  1. Write snapshot 1 with watermark 1000 and create keep1 on snapshot 1.
  2. Write additional data in snapshot 2 with the same watermark 1000.
  3. Expire snapshot 1 while retaining keep1.
  4. Call sys.create_tag_from_watermark(..., watermark => 1000).

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: FileStoreCommitImpl inherits the previous value for null watermarks and takes the maximum for new values.

Please break equal-watermark ties by snapshot ID, preferring the earlier snapshot. Simply changing < to <= would choose a later tagged snapshot in the reverse arrangement, so please cover both directions. The existing Flink implementation shares this comparison; this finding concerns the new Spark API and its stated contract, rather than a Flink regression.

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 (expected snapshot=1, actual=2); independent review confirmed the reachable watermark/tags ordering semantics.

Fixed. When watermarks are equal, the procedure now keeps the earlier snapshot by comparing snapshot IDs, instead of changing < to <=. Using <= would replace an earlier live snapshot with a later tag.

Two tests cover both directions:

  • an expired snapshot retained by a tag is chosen over a later live snapshot with the same watermark
  • a later tag with the same watermark does not replace an earlier live snapshot

This change is limited to the new Spark procedure. The Flink implementation still uses the original comparison.

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.

2 participants