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 fa3c096700 [spark] Fix missing data for delete hiding remaining 
records (#9224)
fa3c096700 is described below

commit fa3c096700f30368c14bf16668c4816c8e9963fe
Author: Arnav Balyan <[email protected]>
AuthorDate: Sat Aug 15 17:44:27 2026 +0530

    [spark] Fix missing data for delete hiding remaining records (#9224)
---
 .../spark/commands/DeleteFromPaimonTableCommand.scala     | 10 ++++++++--
 .../apache/paimon/spark/commands/PaimonSparkWriter.scala  | 12 +++++++++---
 .../org/apache/paimon/spark/write/PaimonDataWrite.scala   |  6 +++++-
 .../apache/paimon/spark/sql/DeleteFromTableTestBase.scala | 15 +++++++++++++++
 4 files changed, 37 insertions(+), 6 deletions(-)

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 247a998e47..e85ef0272e 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,6 +18,7 @@
 
 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
@@ -104,8 +105,13 @@ case class DeleteFromPaimonTableCommand(
         data = selectWithRowTracking(data)
       }
 
-      // only write new files, should have no compaction
-      val addCommitMessage = writer.writeOnly().withRowTracking().write(data)
+      val rewriteWriter =
+        if (coreOptions.mergeEngine() == FIRST_ROW) {
+          writer.withIgnorePreviousFiles()
+        } else {
+          writer.writeOnly()
+        }
+      val addCommitMessage = rewriteWriter.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 b201b88110..35d7098529 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,7 +57,8 @@ import scala.collection.JavaConverters._
 case class PaimonSparkWriter(
     table: FileStoreTable,
     writeRowTracking: Boolean = false,
-    batchId: Option[Long] = None)
+    batchId: Option[Long] = None,
+    ignorePreviousFiles: Boolean = false)
   extends WriteHelper {
 
   private lazy val tableSchema = table.schema
@@ -115,9 +116,13 @@ 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)
+      PaimonSparkWriter(table, writeRowTracking = true, ignorePreviousFiles = 
ignorePreviousFiles)
     } else {
       this
     }
@@ -175,7 +180,8 @@ case class PaimonSparkWriter(
         fullCompactionDeltaCommits,
         batchId,
         uriReaderFactory,
-        postponePartitionBucketComputer
+        postponePartitionBucketComputer,
+        ignorePreviousFiles
       )
 
     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 af20144f52..fd35b6cfa7 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,7 +38,8 @@ case class PaimonDataWrite(
     fullCompactionDeltaCommits: Option[Int],
     batchId: Option[Long],
     uriReaderFactory: UriReaderFactory,
-    postponePartitionBucketComputer: Option[BinaryRow => Integer])
+    postponePartitionBucketComputer: Option[BinaryRow => Integer],
+    ignorePreviousFiles: Boolean = false)
   extends abstractInnerTableDataWrite[Row]
   with InnerTableV1DataWrite {
 
@@ -47,6 +48,9 @@ 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 2b89afcc68..90bddb7aae 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
@@ -252,6 +252,21 @@ abstract class DeleteFromTableTestBase extends 
PaimonSparkTestBase {
       }
   }
 
+  test("Paimon Delete: first-row table") {
+    withTable("t") {
+      sql("""CREATE TABLE t (id INT, name STRING)
+            |TBLPROPERTIES ('primary-key' = 'id', 'bucket' = '1', 
'merge-engine' = 'first-row')
+            |""".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")))
+    }
+  }
+
   test(s"test delete with primary key") {
     spark.sql(
       s"""

Reply via email to