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 7e9403ee09 [spark] supports deleting rows with DELETE operation on
data evolution tables (#8418)
7e9403ee09 is described below
commit 7e9403ee0903aea2e9add6f5fe1f7df608d475c1
Author: Faiz <[email protected]>
AuthorDate: Wed Jul 1 22:57:31 2026 +0800
[spark] supports deleting rows with DELETE operation on data evolution
tables (#8418)
---
.../catalyst/analysis/PaimonDeleteTable.scala | 7 +-
.../catalyst/analysis/PaimonDeleteTable.scala | 7 +-
.../commands/DeleteFromPaimonTableCommand.scala | 3 +-
.../spark/commands/PaimonRowLevelCommand.scala | 133 +++++++++++++++++++--
.../paimon/spark/sql/DeleteFromTableTestBase.scala | 46 +++++++
.../paimon/spark/sql/RowTrackingTestBase.scala | 3 +-
6 files changed, 186 insertions(+), 13 deletions(-)
diff --git
a/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
b/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
index 329c4c5886..707b24824d 100644
---
a/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
+++
b/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
@@ -39,9 +39,12 @@ object PaimonDeleteTable extends Rule[LogicalPlan] with
RowLevelHelper {
table.getTable match {
case paimonTable: FileStoreTable =>
val relation = PaimonRelation.getPaimonRelation(d.table)
- if (paimonTable.coreOptions().dataEvolutionEnabled()) {
+ if (
+ paimonTable.coreOptions().dataEvolutionEnabled()
+ && !paimonTable.coreOptions().deletionVectorsEnabled()
+ ) {
throw new RuntimeException(
- "Delete operation is not supported when data evolution is
enabled yet.")
+ "Can only perform deletion operation on data evolution tables
with DeletionVector enabled.")
}
DeleteFromPaimonTableCommand(relation, paimonTable,
condition.getOrElse(TrueLiteral))
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
index eddf218b26..694b41fba4 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonDeleteTable.scala
@@ -45,9 +45,12 @@ object PaimonDeleteTable extends Rule[LogicalPlan] with
RowLevelHelper {
table.getTable match {
case paimonTable: FileStoreTable =>
val relation = PaimonRelation.getPaimonRelation(d.table)
- if (paimonTable.coreOptions().dataEvolutionEnabled()) {
+ if (
+ paimonTable.coreOptions().dataEvolutionEnabled()
+ && !paimonTable.coreOptions().deletionVectorsEnabled()
+ ) {
throw new RuntimeException(
- "Delete operation is not supported when data evolution is
enabled yet.")
+ "Can only perform deletion operation on data evolution
tables with DeletionVector enabled.")
}
DeleteFromPaimonTableCommand(relation, paimonTable, condition)
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 cd29efa152..247a998e47 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
@@ -79,7 +79,8 @@ case class DeleteFromPaimonTableCommand(
dataFilePathToMeta,
condition,
relation,
- sparkSession)
+ sparkSession,
+ coreOptions.dataEvolutionEnabled())
// Step3: update the touched deletion vectors and index files
writer.persistDeletionVectors(deletionVectors, readSnapshot)
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonRowLevelCommand.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonRowLevelCommand.scala
index f9dc3b14c7..3fad7fb4df 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonRowLevelCommand.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonRowLevelCommand.scala
@@ -22,6 +22,7 @@ import org.apache.paimon.CoreOptions
import org.apache.paimon.deletionvectors.{Bitmap64DeletionVector,
BitmapDeletionVector, DeletionVector}
import org.apache.paimon.fs.Path
import org.apache.paimon.io.{CompactIncrement, DataFileMeta, DataIncrement}
+import org.apache.paimon.operation.DataEvolutionSplitRead
import org.apache.paimon.spark.catalyst.analysis.expressions.ExpressionHelper
import
org.apache.paimon.spark.commands.SparkDataFileMeta.convertToSparkDataFileMeta
import org.apache.paimon.spark.leafnode.PaimonLeafRunnableCommand
@@ -30,6 +31,7 @@ import org.apache.paimon.spark.util.ScanPlanHelper
import org.apache.paimon.table.BucketMode
import org.apache.paimon.table.sink.{CommitMessage, CommitMessageImpl}
import org.apache.paimon.table.source.DataSplit
+import org.apache.paimon.utils.DataEvolutionUtils.retrieveAnchorFile
import org.apache.paimon.utils.SerializationUtils
import org.apache.spark.sql.{Dataset, Row, SparkSession}
@@ -42,6 +44,7 @@ import
org.apache.spark.sql.execution.datasources.v2.DataSourceV2Relation
import org.apache.spark.sql.functions.col
import java.util.Collections
+import java.util.function.Function
import scala.collection.JavaConverters._
@@ -144,17 +147,72 @@ trait PaimonRowLevelCommand
dataFilePathToMeta: Map[String, SparkDataFileMeta],
condition: Expression,
relation: DataSourceV2Relation,
- sparkSession: SparkSession): Dataset[SparkDeletionVector] = {
+ sparkSession: SparkSession,
+ dataEvolutionEnabled: Boolean = false): Dataset[SparkDeletionVector] = {
val filteredRelation =
createNewScanPlan(candidateDataSplits, relation, Some(condition))
val dataset = createDataset(sparkSession, filteredRelation)
- collectDeletionVectors(dataFilePathToMeta, dataset, sparkSession)
+ val deletionTargets =
+ if (dataEvolutionEnabled) {
+ dataEvolutionDeletionTargets(candidateDataSplits, dataset,
sparkSession)
+ } else {
+ pathAndIndexDeletionTargets(dataset, sparkSession)
+ }
+ buildDeletionVectors(dataFilePathToMeta, deletionTargets, sparkSession)
+ }
+
+ /**
+ * Collect deletion targets for data evolution tables: find all [AnchorFile,
Deletion position]
+ * pairs. For deletion positions, we simply use `recordRowId -
anchorFileRangeStart` to get the
+ * offset, instead of calculating returned positions.
+ */
+ private def dataEvolutionDeletionTargets(
+ candidateDataSplits: Seq[DataSplit],
+ dataset: Dataset[Row],
+ sparkSession: SparkSession):
Dataset[PaimonRowLevelCommand.DeletionTarget] = {
+ import sparkSession.implicits._
+
+ val anchorRanges = dataEvolutionAnchorRanges(candidateDataSplits)
+ val broadcastAnchorRanges =
sparkSession.sparkContext.broadcast(anchorRanges)
+
+ dataset
+ .select(col(ROW_ID_COLUMN))
+ .as[Long]
+ .map {
+ rowId =>
+ val (filePath, rowIndex) =
+ PaimonRowLevelCommand.rowIdToPathAndIndex(rowId,
broadcastAnchorRanges.value)
+ PaimonRowLevelCommand.DeletionTarget(filePath, rowIndex)
+ }
}
protected def collectDeletionVectors(
dataFilePathToMeta: Map[String, SparkDataFileMeta],
dataset: Dataset[Row],
sparkSession: SparkSession): Dataset[SparkDeletionVector] = {
+ buildDeletionVectors(
+ dataFilePathToMeta,
+ pathAndIndexDeletionTargets(dataset, sparkSession),
+ sparkSession)
+ }
+
+ private def pathAndIndexDeletionTargets(
+ dataset: Dataset[Row],
+ sparkSession: SparkSession):
Dataset[PaimonRowLevelCommand.DeletionTarget] = {
+ import sparkSession.implicits._
+ dataset
+ .select(PATH_AND_INDEX_META_COLUMNS.map(col): _*)
+ .as[(String, Long)]
+ .map {
+ case (filePath, rowIndex) =>
+ PaimonRowLevelCommand.DeletionTarget(filePath, rowIndex)
+ }
+ }
+
+ private def buildDeletionVectors(
+ dataFilePathToMeta: Map[String, SparkDataFileMeta],
+ deletionTargets: Dataset[PaimonRowLevelCommand.DeletionTarget],
+ sparkSession: SparkSession): Dataset[SparkDeletionVector] = {
import sparkSession.implicits._
// convert to a serializable map
val dataFileNameToPartitionAndBucket = dataFilePathToMeta.map {
@@ -164,16 +222,14 @@ trait PaimonRowLevelCommand
val my_table = table
val dvBitmap64 = my_table.coreOptions().deletionVectorBitmap64()
- dataset
- .select(PATH_AND_INDEX_META_COLUMNS.map(col): _*)
- .as[(String, Long)]
- .groupByKey(_._1)
+ deletionTargets
+ .groupByKey(_.filePath)
.mapGroups {
(filePath, iter) =>
val dv =
if (dvBitmap64) new Bitmap64DeletionVector() else new
BitmapDeletionVector()
while (iter.hasNext) {
- dv.delete(iter.next()._2)
+ dv.delete(iter.next().rowIndex)
}
val (bucketPath, partition, bucket, dataFilePath) =
@@ -188,6 +244,33 @@ trait PaimonRowLevelCommand
}
}
+ private def dataEvolutionAnchorRanges(
+ candidateDataSplits: Seq[DataSplit]):
Array[PaimonRowLevelCommand.AnchorRange] = {
+ val identity =
+ new Function[DataFileMeta, DataFileMeta] {
+ override def apply(file: DataFileMeta): DataFileMeta = file
+ }
+
+ candidateDataSplits
+ .flatMap {
+ split =>
+
DataEvolutionSplitRead.mergeRangesAndSort(split.dataFiles()).asScala.map {
+ group =>
+ val anchor = retrieveAnchorFile(group, identity)
+ val range = anchor.nonNullRowIdRange()
+ val anchorFilePath =
+ if (anchor.externalPath().isPresent) {
+ anchor.externalPath().get()
+ } else {
+ split.bucketPath() + "/" + anchor.fileName()
+ }
+ PaimonRowLevelCommand.AnchorRange(range.from, range.to,
anchorFilePath)
+ }
+ }
+ .sortBy(_.from)
+ .toArray
+ }
+
protected def buildDeletedCommitMessage(
deletedFiles: Array[SparkDataFileMeta]): Seq[CommitMessage] = {
deletedFiles
@@ -213,3 +296,39 @@ trait PaimonRowLevelCommand
.toSeq
}
}
+
+object PaimonRowLevelCommand {
+
+ /** A data-evolution anchor file and its covered row-id range. */
+ final private[commands] case class AnchorRange(from: Long, to: Long,
filePath: String)
+
+ /** A deletion target represented by data file path and local row index. */
+ final private[commands] case class DeletionTarget(filePath: String,
rowIndex: Long)
+
+ /** Binary search to lookup anchor file & range for each row id. */
+ private[commands] def rowIdToPathAndIndex(
+ rowId: Long,
+ anchorRanges: Array[AnchorRange]): (String, Long) = {
+ var low = 0
+ var high = anchorRanges.length - 1
+ var candidate = -1
+
+ while (low <= high) {
+ val mid = (low + high) >>> 1
+ if (anchorRanges(mid).from <= rowId) {
+ candidate = mid
+ low = mid + 1
+ } else {
+ high = mid - 1
+ }
+ }
+
+ if (candidate < 0 || rowId > anchorRanges(candidate).to) {
+ throw new IllegalStateException(
+ s"Cannot find data-evolution deletion-vector anchor range for row id
$rowId.")
+ }
+
+ val anchor = anchorRanges(candidate)
+ (anchor.filePath, rowId - anchor.from)
+ }
+}
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..f45a9463b6 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
@@ -519,6 +519,52 @@ abstract class DeleteFromTableTestBase extends
PaimonSparkTestBase {
}
}
+ test("Paimon delete: data evolution table with deletion vectors") {
+ withTable("t") {
+ sql("""
+ |CREATE TABLE t (id INT, b INT, c INT)
+ |TBLPROPERTIES (
+ | 'row-tracking.enabled' = 'true',
+ | 'data-evolution.enabled' = 'true',
+ | 'deletion-vectors.enabled' = 'true')
+ |""".stripMargin)
+ sql("INSERT INTO t SELECT /*+ REPARTITION(1) */ id, id AS b, id AS c
FROM range(0, 5)")
+ sql("INSERT INTO t SELECT /*+ REPARTITION(1) */ id, id AS b, id AS c
FROM range(5, 10)")
+ sql("ALTER TABLE t ADD COLUMNS (d INT)")
+ sql(
+ "INSERT INTO t SELECT /*+ REPARTITION(1) */ id, id AS b, id AS c, id +
100 AS d FROM range(10, 13)")
+
+ sql("DELETE FROM t WHERE id IN (1, 4, 6, 11)")
+ checkAnswer(
+ sql("SELECT *, _ROW_ID FROM t ORDER BY id"),
+ Seq(
+ Row(0, 0, 0, null, 0L),
+ Row(2, 2, 2, null, 2L),
+ Row(3, 3, 3, null, 3L),
+ Row(5, 5, 5, null, 5L),
+ Row(7, 7, 7, null, 7L),
+ Row(8, 8, 8, null, 8L),
+ Row(9, 9, 9, null, 9L),
+ Row(10, 10, 10, 110, 10L),
+ Row(12, 12, 12, 112, 12L)
+ )
+ )
+
+ sql("DELETE FROM t WHERE id IN (2, 8)")
+ checkAnswer(
+ sql("SELECT *, _ROW_ID FROM t ORDER BY id"),
+ Seq(
+ Row(0, 0, 0, null, 0L),
+ Row(3, 3, 3, null, 3L),
+ Row(5, 5, 5, null, 5L),
+ Row(7, 7, 7, null, 7L),
+ Row(9, 9, 9, null, 9L),
+ Row(10, 10, 10, 110, 10L),
+ Row(12, 12, 12, 112, 12L))
+ )
+ }
+ }
+
test("Paimon delete: delete with range condition") {
withTable("t") {
sql(s"CREATE TABLE t (id INT, v INT)")
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
index f960b59d9b..6c63e1919d 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
@@ -1130,7 +1130,8 @@ abstract class RowTrackingTestBase extends
PaimonSparkTestBase with AdaptiveSpar
intercept[RuntimeException] {
sql("DELETE FROM t WHERE id = 2")
}.getMessage
- .contains("Delete operation is not supported when data evolution is
enabled yet."))
+ .contains(
+ "Can only perform deletion operation on data evolution tables with
DeletionVector enabled."))
}
}