This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 3a4bc9c71d [spark] Forbid delete and update on first-row table (#6718) 
(#9355)
3a4bc9c71d is described below

commit 3a4bc9c71d3393da9a961db81869097e3b6080e9
Author: Lihan Guo <[email protected]>
AuthorDate: Sun Aug 23 23:19:44 2026 +0800

    [spark] Forbid delete and update on first-row table (#6718) (#9355)
---
 .../spark/catalyst/analysis/RowLevelOp.scala       |  7 +---
 .../commands/DeleteFromPaimonTableCommand.scala    |  9 +----
 .../paimon/spark/commands/PaimonSparkWriter.scala  | 12 ++----
 .../paimon/spark/write/PaimonDataWrite.scala       |  6 +--
 .../paimon/spark/sql/DeleteFromTableTestBase.scala | 41 ++++++++++++-------
 .../paimon/spark/sql/MergeIntoTableTestBase.scala  | 47 +++++++++++++++++++++-
 6 files changed, 79 insertions(+), 43 deletions(-)

diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/RowLevelOp.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/RowLevelOp.scala
index 3e1e2b52d2..16803e0f09 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/RowLevelOp.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/RowLevelOp.scala
@@ -46,11 +46,8 @@ sealed trait RowLevelOp {
 
 case object Delete extends RowLevelOp {
 
-  override val supportedMergeEngine: Seq[MergeEngine] = Seq(
-    MergeEngine.DEDUPLICATE,
-    MergeEngine.PARTIAL_UPDATE,
-    MergeEngine.AGGREGATE,
-    MergeEngine.FIRST_ROW)
+  override val supportedMergeEngine: Seq[MergeEngine] =
+    Seq(MergeEngine.DEDUPLICATE, MergeEngine.PARTIAL_UPDATE, 
MergeEngine.AGGREGATE)
 
   override val supportAppendOnlyTable: Boolean = true
 
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DeleteFromPaimonTableCommand.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DeleteFromPaimonTableCommand.scala
index e85ef0272e..b57dc9cb1a 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DeleteFromPaimonTableCommand.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DeleteFromPaimonTableCommand.scala
@@ -18,7 +18,6 @@
 
 package org.apache.paimon.spark.commands
 
-import org.apache.paimon.CoreOptions.MergeEngine.FIRST_ROW
 import org.apache.paimon.Snapshot
 import org.apache.paimon.spark.catalyst.analysis.expressions.ExpressionHelper
 import org.apache.paimon.spark.schema.SparkSystemColumns.ROW_KIND_COL
@@ -105,13 +104,7 @@ case class DeleteFromPaimonTableCommand(
         data = selectWithRowTracking(data)
       }
 
-      val rewriteWriter =
-        if (coreOptions.mergeEngine() == FIRST_ROW) {
-          writer.withIgnorePreviousFiles()
-        } else {
-          writer.writeOnly()
-        }
-      val addCommitMessage = rewriteWriter.withRowTracking().write(data)
+      val addCommitMessage = writer.writeOnly().withRowTracking().write(data)
 
       // Step5: convert the deleted files that need to be written to commit 
message.
       val deletedCommitMessage = buildDeletedCommitMessage(touchedFiles)
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
index 2bbd53b8ab..91efc3d541 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
@@ -57,8 +57,7 @@ import scala.collection.JavaConverters._
 case class PaimonSparkWriter(
     table: FileStoreTable,
     writeRowTracking: Boolean = false,
-    batchId: Option[Long] = None,
-    ignorePreviousFiles: Boolean = false)
+    batchId: Option[Long] = None)
   extends WriteHelper {
 
   private lazy val tableSchema = table.schema
@@ -116,13 +115,9 @@ case class PaimonSparkWriter(
     PaimonSparkWriter(table.copy(singletonMap(WRITE_ONLY.key(), "true")))
   }
 
-  def withIgnorePreviousFiles(): PaimonSparkWriter = {
-    copy(ignorePreviousFiles = true)
-  }
-
   def withRowTracking(): PaimonSparkWriter = {
     if (coreOptions.rowTrackingEnabled()) {
-      PaimonSparkWriter(table, writeRowTracking = true, ignorePreviousFiles = 
ignorePreviousFiles)
+      PaimonSparkWriter(table, writeRowTracking = true)
     } else {
       this
     }
@@ -180,8 +175,7 @@ case class PaimonSparkWriter(
         fullCompactionDeltaCommits,
         batchId,
         uriReaderFactory,
-        postponePartitionBucketComputer,
-        ignorePreviousFiles
+        postponePartitionBucketComputer
       )
 
     def sparkParallelism = {
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
index fd35b6cfa7..af20144f52 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
@@ -38,8 +38,7 @@ case class PaimonDataWrite(
     fullCompactionDeltaCommits: Option[Int],
     batchId: Option[Long],
     uriReaderFactory: UriReaderFactory,
-    postponePartitionBucketComputer: Option[BinaryRow => Integer],
-    ignorePreviousFiles: Boolean = false)
+    postponePartitionBucketComputer: Option[BinaryRow => Integer])
   extends abstractInnerTableDataWrite[Row]
   with InnerTableV1DataWrite {
 
@@ -48,9 +47,6 @@ case class PaimonDataWrite(
   val write: TableWriteImpl[Row] = {
     val _write = writeBuilder.newWrite().asInstanceOf[TableWriteImpl[Row]]
     _write.withIOManager(ioManager)
-    if (ignorePreviousFiles) {
-      _write.withIgnorePreviousFiles(true)
-    }
     if (writeRowTracking) {
       _write.withWriteType(writeType)
     }
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DeleteFromTableTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DeleteFromTableTestBase.scala
index 90bddb7aae..5369e6e6a2 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DeleteFromTableTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DeleteFromTableTestBase.scala
@@ -21,6 +21,7 @@ package org.apache.paimon.spark.sql
 import org.apache.paimon.{CoreOptions, Snapshot}
 import org.apache.paimon.CoreOptions.MergeEngine
 import org.apache.paimon.spark.PaimonSparkTestBase
+import org.apache.paimon.spark.catalyst.analysis.Delete
 
 import org.apache.spark.sql.Row
 import org.assertj.core.api.Assertions.{assertThat, assertThatThrownBy}
@@ -239,15 +240,21 @@ abstract class DeleteFromTableTestBase extends 
PaimonSparkTestBase {
           spark.sql("INSERT INTO T VALUES (2, 'b', NULL)")
           spark.sql("INSERT INTO T VALUES (1, NULL, 16)")
 
-          if (mergeEngine != MergeEngine.DEDUPLICATE) {
+          if (Delete.supportedMergeEngine.contains(mergeEngine)) {
+            if (mergeEngine != MergeEngine.DEDUPLICATE) {
+              assertThatThrownBy(() => spark.sql("DELETE FROM T WHERE id = 1"))
+                .hasMessageContaining("please use 'COMPACT' procedure first")
+              spark.sql("CALL sys.compact(table => 'T')")
+            }
+
+            spark.sql("DELETE FROM T WHERE id = 1")
+            assertThat(spark.sql("SELECT * FROM T").collectAsList().toString)
+              .isEqualTo("[[2,b,null]]")
+          } else {
             assertThatThrownBy(() => spark.sql("DELETE FROM T WHERE id = 1"))
-              .hasMessageContaining("please use 'COMPACT' procedure first")
-            spark.sql("CALL sys.compact(table => 'T')")
+              .isInstanceOf(classOf[UnsupportedOperationException])
+              .hasMessageContaining(s"merge engine $mergeEngine can not 
support Delete")
           }
-
-          spark.sql("DELETE FROM T WHERE id = 1")
-          assertThat(spark.sql("SELECT * FROM T").collectAsList().toString)
-            .isEqualTo("[[2,b,null]]")
         }
       }
   }
@@ -259,11 +266,9 @@ abstract class DeleteFromTableTestBase extends 
PaimonSparkTestBase {
             |""".stripMargin)
       sql("INSERT INTO t VALUES (1, 'a'), (2, 'b'), (3, 'c')")
 
-      checkAnswer(sql("SELECT * FROM t ORDER BY id"), Seq(Row(1, "a"), Row(2, 
"b"), Row(3, "c")))
-
-      sql("DELETE FROM t WHERE id = 3")
-
-      checkAnswer(sql("SELECT * FROM t ORDER BY id"), Seq(Row(1, "a"), Row(2, 
"b")))
+      assertThatThrownBy(() => sql("DELETE FROM t WHERE id = 3"))
+        .isInstanceOf(classOf[UnsupportedOperationException])
+        .hasMessageContaining("merge engine first-row can not support Delete")
     }
   }
 
@@ -491,9 +496,15 @@ abstract class DeleteFromTableTestBase extends 
PaimonSparkTestBase {
           // update
           spark.sql("INSERT INTO T VALUES (1, NULL, 16)")
           // delete
-          spark.sql("DELETE FROM T WHERE id = 1")
-          assertThat(spark.sql("SELECT * FROM T").collectAsList().toString)
-            .isEqualTo("[[2,b,null]]")
+          if (Delete.supportedMergeEngine.contains(mergeEngine)) {
+            spark.sql("DELETE FROM T WHERE id = 1")
+            assertThat(spark.sql("SELECT * FROM T").collectAsList().toString)
+              .isEqualTo("[[2,b,null]]")
+          } else {
+            assertThatThrownBy(() => spark.sql("DELETE FROM T WHERE id = 1"))
+              .isInstanceOf(classOf[UnsupportedOperationException])
+              .hasMessageContaining(s"merge engine $mergeEngine can not 
support Delete")
+          }
         }
       }
   }
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoTableTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoTableTestBase.scala
index 78039e6083..014c953b41 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoTableTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoTableTestBase.scala
@@ -18,10 +18,12 @@
 
 package org.apache.paimon.spark.sql
 
-import org.apache.paimon.Snapshot
+import org.apache.paimon.{CoreOptions, Snapshot}
 import org.apache.paimon.spark.{PaimonAppendTable, PaimonPrimaryKeyTable, 
PaimonSparkTestBase, PaimonTableTest}
+import org.apache.paimon.spark.catalyst.analysis.MergeInto
 
 import org.apache.spark.sql.{AnalysisException, Row}
+import org.assertj.core.api.Assertions.{assertThat, assertThatThrownBy}
 
 import java.util.UUID
 import java.util.concurrent.Executors
@@ -1832,6 +1834,49 @@ trait MergeIntoAppendTableTest extends 
PaimonSparkTestBase with PaimonAppendTabl
     }
   }
 
+  CoreOptions.MergeEngine.values().foreach {
+    mergeEngine =>
+      {
+        test(s"test merge into with merge engine $mergeEngine") {
+          val options = if ("first-row".equals(mergeEngine.toString)) {
+            s"'primary-key' = 'id', 'merge-engine' = '$mergeEngine', 
'changelog-producer' = 'lookup'"
+          } else {
+            s"'primary-key' = 'id', 'merge-engine' = '$mergeEngine'"
+          }
+          withTable("source", "target") {
+            Seq((1, "a_new", "11")).toDF("id", "name", 
"dt").createOrReplaceTempView("source")
+            spark.sql(s"""
+                         |CREATE TABLE target (id INT, name STRING, dt STRING)
+                         |TBLPROPERTIES ($options)
+                         |""".stripMargin)
+            spark.sql("INSERT INTO target VALUES (1, 'a', '11'), (2, 'b', 
'22')")
+
+            if (MergeInto.supportedMergeEngine.contains(mergeEngine)) {
+              spark.sql("""
+                          |MERGE INTO target
+                          |USING source
+                          |ON target.id = source.id
+                          |WHEN MATCHED THEN
+                          |UPDATE SET target.name = source.name
+                          |""".stripMargin)
+              val rows = spark.sql("SELECT * FROM target ORDER BY 
id").collectAsList()
+              assertThat(rows.toString).isEqualTo("[[1,a_new,11], [2,b,22]]")
+            } else {
+              assertThatThrownBy(() => spark.sql("""
+                                                   |MERGE INTO target
+                                                   |USING source
+                                                   |ON target.id = source.id
+                                                   |WHEN MATCHED THEN
+                                                   |UPDATE SET target.name = 
source.name
+                                                   |""".stripMargin))
+                .isInstanceOf(classOf[UnsupportedOperationException])
+                .hasMessageContaining(s"merge engine $mergeEngine can not 
support MergeInto")
+            }
+          }
+        }
+      }
+  }
+
   def createPositiveRandomInt(): Int = {
     val random = new Random()
     val positiveInt = random.nextInt()

Reply via email to