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"""