From 6f4c0b585d09f5b1358f3f9adae1a03ed44539c7 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Sat, 26 Sep 2026 17:17:04 +0800 Subject: [PATCH 1/4] [spark] Add create_tag_from_watermark procedure Flink has a create_tag_from_watermark procedure but Spark has none, so a Spark user cannot tag the first snapshot whose watermark reaches a target -- even though Spark already exposes rollback_to_watermark and create_tag_from_timestamp. It is useful when a Flink writer commits event-time watermarks and a Spark job needs a consistent tag at a watermark boundary. Add CreateTagFromWatermarkProcedure, registered as create_tag_from_watermark. It mirrors the existing CreateTagFromTimestampProcedure structure and the Flink procedure's logic: resolve SnapshotManager.laterOrEqualWatermark, fall back to the earliest tag whose watermark already covers the target, else raise SnapshotNotExistException. Tests: CreateTagFromWatermarkProcedureTest commits three snapshots with watermarks 1000/2000/3000 (through the stream write builder, since a Spark batch write carries no watermark) and asserts the resolved snapshot for a below/equal watermark and a SnapshotNotExistException past the last one. Written with Claude Code; verification is mine. --- .../apache/paimon/spark/SparkProcedures.java | 3 + .../CreateTagFromWatermarkProcedure.java | 139 ++++++++++++++++++ .../CreateTagFromWatermarkProcedureTest.scala | 89 +++++++++++ 3 files changed, 231 insertions(+) create mode 100644 paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java create mode 100644 paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java index f534773959f6..fb6a3c805813 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java @@ -31,6 +31,7 @@ import org.apache.paimon.spark.procedure.CreateGlobalIndexProcedure; import org.apache.paimon.spark.procedure.CreatePolicyProcedure; import org.apache.paimon.spark.procedure.CreateTagFromTimestampProcedure; +import org.apache.paimon.spark.procedure.CreateTagFromWatermarkProcedure; import org.apache.paimon.spark.procedure.CreateTagProcedure; import org.apache.paimon.spark.procedure.DeleteBranchProcedure; import org.apache.paimon.spark.procedure.DeleteTagProcedure; @@ -106,6 +107,8 @@ private static Map> initProcedureBuilders() { procedureBuilders.put("rename_tag", RenameTagProcedure::builder); procedureBuilders.put( "create_tag_from_timestamp", CreateTagFromTimestampProcedure::builder); + procedureBuilders.put( + "create_tag_from_watermark", CreateTagFromWatermarkProcedure::builder); procedureBuilders.put("delete_tag", DeleteTagProcedure::builder); procedureBuilders.put("expire_tags", ExpireTagsProcedure::builder); procedureBuilders.put("create_branch", CreateBranchProcedure::builder); diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java new file mode 100644 index 000000000000..5336c3fcd77d --- /dev/null +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java @@ -0,0 +1,139 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.spark.procedure; + +import org.apache.paimon.Snapshot; +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.utils.SnapshotManager; +import org.apache.paimon.utils.SnapshotNotExistException; +import org.apache.paimon.utils.TimeUtils; + +import org.apache.spark.sql.catalyst.InternalRow; +import org.apache.spark.sql.connector.catalog.Identifier; +import org.apache.spark.sql.connector.catalog.TableCatalog; +import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.Metadata; +import org.apache.spark.sql.types.StructField; +import org.apache.spark.sql.types.StructType; +import org.apache.spark.unsafe.types.UTF8String; + +import java.time.Duration; +import java.util.Set; + +import static org.apache.spark.sql.types.DataTypes.LongType; +import static org.apache.spark.sql.types.DataTypes.StringType; + +/** The procedure supports creating tags from a snapshot's watermark. */ +public class CreateTagFromWatermarkProcedure extends BaseProcedure { + + private static final ProcedureParameter[] PARAMETERS = + new ProcedureParameter[] { + ProcedureParameter.required("table", StringType), + ProcedureParameter.required("tag", StringType), + ProcedureParameter.required("watermark", LongType), + ProcedureParameter.optional("time_retained", StringType) + }; + + private static final StructType OUTPUT_TYPE = + new StructType( + new StructField[] { + new StructField("tagName", DataTypes.StringType, true, Metadata.empty()), + new StructField("snapshot", DataTypes.LongType, true, Metadata.empty()), + new StructField("commit_time", DataTypes.LongType, true, Metadata.empty()), + new StructField("watermark", DataTypes.StringType, true, Metadata.empty()) + }); + + protected CreateTagFromWatermarkProcedure(TableCatalog tableCatalog) { + super(tableCatalog); + } + + @Override + public ProcedureParameter[] parameters() { + return PARAMETERS; + } + + @Override + public StructType outputType() { + return OUTPUT_TYPE; + } + + @Override + public InternalRow[] call(InternalRow args) { + Identifier tableIdent = toIdentifier(args.getString(0), PARAMETERS[0].name()); + String tag = args.getString(1); + long watermark = args.getLong(2); + Duration timeRetained = + args.isNullAt(3) ? null : TimeUtils.parseDuration(args.getString(3)); + + return modifyPaimonTable( + tableIdent, + table -> { + FileStoreTable fileStoreTable = (FileStoreTable) table; + SnapshotManager snapshotManager = fileStoreTable.snapshotManager(); + Snapshot snapshot = snapshotManager.laterOrEqualWatermark(watermark); + + Set sortedTagsSnapshots = fileStoreTable.tagManager().tags().keySet(); + + // Fall back to the earliest tag whose watermark already covers the target, + // mirroring the Flink create_tag_from_watermark procedure. + for (Snapshot tagSnapshot : sortedTagsSnapshots) { + if (tagSnapshot.watermark() != null + && watermark <= tagSnapshot.watermark()) { + if (snapshot == null + || snapshot.watermark() == null + || tagSnapshot.watermark() < snapshot.watermark()) { + snapshot = tagSnapshot; + } + break; + } + } + + SnapshotNotExistException.checkNotNull( + snapshot, + String.format( + "Could not find any snapshot whose watermark later than %s.", + watermark)); + + fileStoreTable.createTag(tag, snapshot.id(), timeRetained); + + InternalRow outputRow = + newInternalRow( + UTF8String.fromString(tag), + snapshot.id(), + snapshot.timeMillis(), + UTF8String.fromString(String.valueOf(snapshot.watermark()))); + + return new InternalRow[] {outputRow}; + }); + } + + public static ProcedureBuilder builder() { + return new BaseProcedure.Builder() { + @Override + public CreateTagFromWatermarkProcedure doBuild() { + return new CreateTagFromWatermarkProcedure(tableCatalog()); + } + }; + } + + @Override + public String description() { + return "CreateTagFromWatermarkProcedure"; + } +} diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala new file mode 100644 index 000000000000..0664a6b94089 --- /dev/null +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala @@ -0,0 +1,89 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.spark.procedure + +import org.apache.paimon.data.{BinaryString, GenericRow} +import org.apache.paimon.manifest.ManifestCommittable +import org.apache.paimon.spark.PaimonSparkTestBase +import org.apache.paimon.table.FileStoreTable +import org.apache.paimon.table.sink.TableCommitImpl +import org.apache.paimon.utils.SnapshotNotExistException + +import org.apache.spark.sql.Row + +import scala.collection.JavaConverters._ + +class CreateTagFromWatermarkProcedureTest extends PaimonSparkTestBase { + + test("Paimon Procedure: create tags from snapshots watermark") { + spark.sql(s""" + |CREATE TABLE T (a INT, b STRING) + |TBLPROPERTIES ('primary-key'='a', 'bucket'='1') + |""".stripMargin) + + val table = loadTable("T") + // A Spark batch write carries no watermark, so commit snapshots with explicit + // watermarks the way a Flink writer would, mirroring the Flink IT case. + writeWithWatermark(table, 1L, 1000L, GenericRow.of(1, BinaryString.fromString("a"))) + writeWithWatermark(table, 2L, 2000L, GenericRow.of(2, BinaryString.fromString("b"))) + writeWithWatermark(table, 3L, 3000L, GenericRow.of(3, BinaryString.fromString("c"))) + + val commitTime1 = table.snapshotManager.snapshot(1).timeMillis + val commitTime2 = table.snapshotManager.snapshot(2).timeMillis + + // watermark below snapshot 1 resolves to snapshot 1. + checkAnswer( + spark.sql(s"""CALL paimon.sys.create_tag_from_watermark( + |table => 'test.T', tag => 'tag1', watermark => 500)""".stripMargin), + Row("tag1", 1, commitTime1, "1000") :: Nil + ) + + // watermark equal to snapshot 2 resolves to snapshot 2. + checkAnswer( + spark.sql(s"""CALL paimon.sys.create_tag_from_watermark( + |table => 'test.T', tag => 'tag2', watermark => 2000)""".stripMargin), + Row("tag2", 2, commitTime2, "2000") :: Nil + ) + + // watermark later than the latest snapshot throws SnapshotNotExistException. + assertThrows[SnapshotNotExistException] { + spark.sql(s"""CALL paimon.sys.create_tag_from_watermark( + |table => 'test.T', tag => 'tag3', watermark => ${Long.MaxValue})""".stripMargin) + } + } + + private def writeWithWatermark( + table: FileStoreTable, + commitId: Long, + watermark: Long, + rows: GenericRow*): Unit = { + val writeBuilder = table.newStreamWriteBuilder() + val writer = writeBuilder.newWrite() + val commit = writeBuilder.newCommit().asInstanceOf[TableCommitImpl] + try { + rows.foreach(writer.write) + val committable = new ManifestCommittable(commitId, watermark) + writer.prepareCommit(true, commitId).asScala.foreach(committable.addFileCommittable) + commit.commit(committable) + } finally { + writer.close() + commit.close() + } + } +} From 002e5f96b1ff994722c4194ae6bf8f9c5d531db0 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Wed, 30 Sep 2026 19:54:04 +0800 Subject: [PATCH 2/4] Document the Spark create_tag_from_watermark procedure --- docs/docs/spark/procedures.md | 1 + docs/docs/spark/procedures/versions.md | 20 ++++++++++++++++++++ 2 files changed, 21 insertions(+) diff --git a/docs/docs/spark/procedures.md b/docs/docs/spark/procedures.md index 2048c49b0718..5ef85a65eb8a 100644 --- a/docs/docs/spark/procedures.md +++ b/docs/docs/spark/procedures.md @@ -78,6 +78,7 @@ Choose a group, then use the page contents to jump to a procedure: [`create_tag`](./procedures/versions#create_tag), [`create_tag_from_timestamp`](./procedures/versions#create_tag_from_timestamp), +[`create_tag_from_watermark`](./procedures/versions#create_tag_from_watermark), [`replace_tag`](./procedures/versions#replace_tag), [`rename_tag`](./procedures/versions#rename_tag), [`delete_tag`](./procedures/versions#delete_tag), diff --git a/docs/docs/spark/procedures/versions.md b/docs/docs/spark/procedures/versions.md index d9b96729722c..7fca7e7f46ef 100644 --- a/docs/docs/spark/procedures/versions.md +++ b/docs/docs/spark/procedures/versions.md @@ -69,6 +69,26 @@ CALL sys.create_tag_from_timestamp( ); ``` +## create_tag_from_watermark + +Create a tag based on given watermark timestamp. + +**Arguments** + +- `table` (`STRING`, required): the target table identifier. +- `tag` (`STRING`, required): name of the new tag. +- `watermark` (`BIGINT`, required): Find the first snapshot (including tagged snapshots) whose watermark is at or after this value. +- `time_retained` (`STRING`, optional): The maximum time retained for newly created tags. + +```sql +CALL sys.create_tag_from_watermark( + `table` => 'default.T', + `tag` => 'my_tag', + `watermark` => 1724404318750, + time_retained => '1 d' +); +``` + ## replace_tag Replace an existing tag with new tag info. From 2b9ebee705941753ec88d7e3d5339ddb7b52a445 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Fri, 2 Oct 2026 22:20:31 +0800 Subject: [PATCH 3/4] [spark] Prefer the earlier snapshot on a create_tag_from_watermark tie --- .../CreateTagFromWatermarkProcedure.java | 10 ++++++- .../CreateTagFromWatermarkProcedureTest.scala | 27 +++++++++++++++++++ 2 files changed, 36 insertions(+), 1 deletion(-) diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java index 5336c3fcd77d..e82b88f0393e 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java @@ -95,9 +95,17 @@ public InternalRow[] call(InternalRow args) { for (Snapshot tagSnapshot : sortedTagsSnapshots) { if (tagSnapshot.watermark() != null && watermark <= tagSnapshot.watermark()) { + long tagWatermark = tagSnapshot.watermark(); + // Prefer the smaller watermark; on a tie prefer the smaller snapshot + // id so the first qualifying snapshot wins. A strict `<` would miss an + // older tagged snapshot that shares the retained candidate's watermark + // (returning a later dataset); a plain `<=` would wrongly let a newer + // equal-watermark tag replace the earlier retained snapshot. if (snapshot == null || snapshot.watermark() == null - || tagSnapshot.watermark() < snapshot.watermark()) { + || tagWatermark < snapshot.watermark() + || (tagWatermark == snapshot.watermark() + && tagSnapshot.id() < snapshot.id())) { snapshot = tagSnapshot; } break; diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala index 0664a6b94089..a2cd9f94469c 100644 --- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala @@ -68,6 +68,33 @@ class CreateTagFromWatermarkProcedureTest extends PaimonSparkTestBase { } } + test("Paimon Procedure: create tag prefers the earlier snapshot on a watermark tie") { + spark.sql(s""" + |CREATE TABLE T (a INT, b STRING) + |TBLPROPERTIES ('primary-key'='a', 'bucket'='1') + |""".stripMargin) + + val table = loadTable("T") + // Two snapshots share watermark 1000; snapshot 1 is tagged, then expired, so + // only the tag preserves it. laterOrEqualWatermark then returns snapshot 2, + // and the tag fallback must still select the earlier snapshot 1 on the tie. + writeWithWatermark(table, 1L, 1000L, GenericRow.of(1, BinaryString.fromString("row1"))) + val commitTime1 = table.snapshotManager.snapshot(1).timeMillis + spark.sql("CALL paimon.sys.create_tag(table => 'test.T', tag => 'historical', snapshot => 1)") + writeWithWatermark(table, 2L, 1000L, GenericRow.of(2, BinaryString.fromString("row2"))) + spark.sql("CALL paimon.sys.expire_snapshots(table => 'test.T', retain_max => 1)") + assert(!table.snapshotManager.snapshotExists(1)) + + // The first qualifying snapshot for watermark 500 is snapshot 1 (smaller id, + // equal watermark), reachable through the tag -- not the retained snapshot 2. + checkAnswer( + spark.sql(s"""CALL paimon.sys.create_tag_from_watermark( + |table => 'test.T', tag => 'boundary', watermark => 500)""".stripMargin), + Row("boundary", 1, commitTime1, "1000") :: Nil + ) + checkAnswer(spark.sql("SELECT * FROM T VERSION AS OF 'boundary'"), Row(1, "row1") :: Nil) + } + private def writeWithWatermark( table: FileStoreTable, commitId: Long, From f97d6567380119721610b9d922d1507ca4b6d109 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Fri, 2 Oct 2026 23:18:48 +0800 Subject: [PATCH 4/4] [spark] Fix expire_snapshots args in the watermark-tie test --- .../spark/procedure/CreateTagFromWatermarkProcedureTest.scala | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala index a2cd9f94469c..3b4e299d1ef8 100644 --- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala @@ -82,7 +82,8 @@ class CreateTagFromWatermarkProcedureTest extends PaimonSparkTestBase { val commitTime1 = table.snapshotManager.snapshot(1).timeMillis spark.sql("CALL paimon.sys.create_tag(table => 'test.T', tag => 'historical', snapshot => 1)") writeWithWatermark(table, 2L, 1000L, GenericRow.of(2, BinaryString.fromString("row2"))) - spark.sql("CALL paimon.sys.expire_snapshots(table => 'test.T', retain_max => 1)") + spark.sql( + "CALL paimon.sys.expire_snapshots(table => 'test.T', retain_max => 1, retain_min => 1)") assert(!table.snapshotManager.snapshotExists(1)) // The first qualifying snapshot for watermark 500 is snapshot 1 (smaller id,