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 475be566fe [spark] Fix batch delete for partial update sequence groups
(#9546)
475be566fe is described below
commit 475be566fef490ad147deec4ee15344c25c0352d
Author: Arnav Balyan <[email protected]>
AuthorDate: Fri Sep 4 12:55:17 2026 +0530
[spark] Fix batch delete for partial update sequence groups (#9546)
---
.../spark/commands/DeleteFromPaimonTableCommand.scala | 5 +++--
.../paimon/spark/sql/DeleteFromTableTestBase.scala | 17 +++++++++++++++++
2 files changed, 20 insertions(+), 2 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 b57dc9cb1a..e8cae8ea52 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,7 @@
package org.apache.paimon.spark.commands
-import org.apache.paimon.Snapshot
+import org.apache.paimon.{CoreOptions, Snapshot}
import org.apache.paimon.spark.catalyst.analysis.expressions.ExpressionHelper
import org.apache.paimon.spark.schema.SparkSystemColumns.ROW_KIND_COL
import org.apache.paimon.table.FileStoreTable
@@ -54,7 +54,8 @@ case class DeleteFromPaimonTableCommand(
private def usePKUpsertDelete(): Boolean = {
try {
validatePKUpsertDeletable(table)
- true
+ coreOptions.mergeEngine() != CoreOptions.MergeEngine.PARTIAL_UPDATE ||
+
coreOptions.toConfiguration.get(CoreOptions.PARTIAL_UPDATE_REMOVE_RECORD_ON_DELETE)
} catch {
case _: UnsupportedOperationException => false
}
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 5369e6e6a2..300c6b7b2d 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
@@ -509,6 +509,23 @@ abstract class DeleteFromTableTestBase extends
PaimonSparkTestBase {
}
}
+ test("Paimon Delete: partial update with remove record on sequence group") {
+ spark.sql(s"""
+ |CREATE TABLE T (id INT, g INT, v BIGINT)
+ |TBLPROPERTIES (
+ | 'primary-key' = 'id',
+ | 'bucket' = '2',
+ | 'merge-engine' = 'partial-update',
+ | 'fields.g.sequence-group' = 'v',
+ | 'partial-update.remove-record-on-sequence-group' = 'g')
+ |""".stripMargin)
+
+ spark.sql("INSERT INTO T VALUES (1, 1, 10), (2, 1, 20)")
+ spark.sql("DELETE FROM T WHERE id = 1")
+
+ checkAnswer(spark.sql("SELECT * FROM T"), Row(2, 1, 20L))
+ }
+
test("Paimon delete: non pk table commit kind") {
for (dvEnabled <- Seq(true, false)) {
withTable("t") {