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. 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..2589aa07c860 --- /dev/null +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedure.java @@ -0,0 +1,177 @@ +/* + * 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 = + earliestSnapshotCoveringWatermark( + snapshotManager, + snapshotManager.laterOrEqualWatermark(watermark), + 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()) { + 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 + || tagWatermark < snapshot.watermark() + || (tagWatermark == snapshot.watermark() + && tagSnapshot.id() < snapshot.id())) { + 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}; + }); + } + + /** + * {@link SnapshotManager#laterOrEqualWatermark} binary-searches and stops on the first + * watermark match it lands on, so a run of snapshots sharing the target watermark can resolve + * to a later member of that run. This Spark API promises the first qualifying snapshot, so walk + * back to the earliest snapshot whose watermark still covers {@code watermark}, stopping at any + * gap or missing-watermark snapshot. + */ + private static Snapshot earliestSnapshotCoveringWatermark( + SnapshotManager snapshotManager, Snapshot snapshot, long watermark) { + if (snapshot == null) { + return null; + } + Long earliest = snapshotManager.earliestSnapshotId(); + while (earliest != null + && snapshot.id() > earliest + && snapshotManager.snapshotExists(snapshot.id() - 1)) { + Snapshot previous = snapshotManager.snapshot(snapshot.id() - 1); + Long previousWatermark = previous.watermark(); + if (previousWatermark == null || previousWatermark < watermark) { + break; + } + snapshot = previous; + } + return snapshot; + } + + 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..42a7238e9be3 --- /dev/null +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CreateTagFromWatermarkProcedureTest.scala @@ -0,0 +1,146 @@ +/* + * 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) + } + } + + 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, retain_min => 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) + } + + test("Paimon Procedure: create tag picks the earliest snapshot in a repeated-watermark run") { + spark.sql(s""" + |CREATE TABLE T (a INT, b STRING) + |TBLPROPERTIES ('primary-key'='a', 'bucket'='1', 'write-only'='true') + |""".stripMargin) + + val table = loadTable("T") + // Three consecutive snapshots (2, 3, 4) share watermark 2000. laterOrEqualWatermark + // binary-searches and stops on the run member it lands on (snapshot 3), but the first + // snapshot covering watermark 2000 is snapshot 2, so the tag must point there. + writeWithWatermark(table, 1L, 1000L, GenericRow.of(1, BinaryString.fromString("row1"))) + writeWithWatermark(table, 2L, 2000L, GenericRow.of(2, BinaryString.fromString("row2"))) + writeWithWatermark(table, 3L, 2000L, GenericRow.of(3, BinaryString.fromString("row3"))) + writeWithWatermark(table, 4L, 2000L, GenericRow.of(4, BinaryString.fromString("row4"))) + writeWithWatermark(table, 5L, 3000L, GenericRow.of(5, BinaryString.fromString("row5"))) + + val commitTime2 = table.snapshotManager.snapshot(2).timeMillis + + checkAnswer( + spark.sql(s"""CALL paimon.sys.create_tag_from_watermark( + |table => 'test.T', tag => 'wm2000', watermark => 2000)""".stripMargin), + Row("wm2000", 2, commitTime2, "2000") :: Nil + ) + checkAnswer( + spark.sql("SELECT * FROM T VERSION AS OF 'wm2000' ORDER BY a"), + Row(1, "row1") :: Row(2, "row2") :: Nil + ) + } + + 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() + } + } +}