Skip to content

Commit 475be56

Browse files
authored
[spark] Fix batch delete for partial update sequence groups (#9546)
1 parent 0289f3d commit 475be56

2 files changed

Lines changed: 20 additions & 2 deletions

File tree

paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DeleteFromPaimonTableCommand.scala

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
package org.apache.paimon.spark.commands
2020

21-
import org.apache.paimon.Snapshot
21+
import org.apache.paimon.{CoreOptions, Snapshot}
2222
import org.apache.paimon.spark.catalyst.analysis.expressions.ExpressionHelper
2323
import org.apache.paimon.spark.schema.SparkSystemColumns.ROW_KIND_COL
2424
import org.apache.paimon.table.FileStoreTable
@@ -54,7 +54,8 @@ case class DeleteFromPaimonTableCommand(
5454
private def usePKUpsertDelete(): Boolean = {
5555
try {
5656
validatePKUpsertDeletable(table)
57-
true
57+
coreOptions.mergeEngine() != CoreOptions.MergeEngine.PARTIAL_UPDATE ||
58+
coreOptions.toConfiguration.get(CoreOptions.PARTIAL_UPDATE_REMOVE_RECORD_ON_DELETE)
5859
} catch {
5960
case _: UnsupportedOperationException => false
6061
}

paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DeleteFromTableTestBase.scala

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -509,6 +509,23 @@ abstract class DeleteFromTableTestBase extends PaimonSparkTestBase {
509509
}
510510
}
511511

512+
test("Paimon Delete: partial update with remove record on sequence group") {
513+
spark.sql(s"""
514+
|CREATE TABLE T (id INT, g INT, v BIGINT)
515+
|TBLPROPERTIES (
516+
| 'primary-key' = 'id',
517+
| 'bucket' = '2',
518+
| 'merge-engine' = 'partial-update',
519+
| 'fields.g.sequence-group' = 'v',
520+
| 'partial-update.remove-record-on-sequence-group' = 'g')
521+
|""".stripMargin)
522+
523+
spark.sql("INSERT INTO T VALUES (1, 1, 10), (2, 1, 20)")
524+
spark.sql("DELETE FROM T WHERE id = 1")
525+
526+
checkAnswer(spark.sql("SELECT * FROM T"), Row(2, 1, 20L))
527+
}
528+
512529
test("Paimon delete: non pk table commit kind") {
513530
for (dvEnabled <- Seq(true, false)) {
514531
withTable("t") {

0 commit comments

Comments
 (0)