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 8a8f9f02df [spark] Use the self-merge shortcut for conditional data 
evolution UPDATE (#10037)
8a8f9f02df is described below

commit 8a8f9f02df2e1cd941b43c2384bfef7ce831d12e
Author: Xiangyi Zhu <[email protected]>
AuthorDate: Mon Sep 21 15:05:24 2026 +0800

    [spark] Use the self-merge shortcut for conditional data evolution UPDATE 
(#10037)
---
 .../UpdatePaimonDataEvolutionTableCommand.scala    |  58 ++++++--
 .../paimon/spark/sql/RowTrackingTestBase.scala     | 153 ++++++++++++++++++++-
 2 files changed, 195 insertions(+), 16 deletions(-)

diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/UpdatePaimonDataEvolutionTableCommand.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/UpdatePaimonDataEvolutionTableCommand.scala
index d3506ecec2..9eac851acf 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/UpdatePaimonDataEvolutionTableCommand.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/UpdatePaimonDataEvolutionTableCommand.scala
@@ -26,7 +26,7 @@ import org.apache.paimon.spark.util.OptionUtils
 
 import org.apache.spark.internal.Logging
 import org.apache.spark.sql.{Row, SparkSession}
-import org.apache.spark.sql.catalyst.expressions.{Alias, Attribute, 
AttributeReference, EqualTo, Expression}
+import org.apache.spark.sql.catalyst.expressions.{Alias, Attribute, 
AttributeReference, EqualTo, Expression, SubqueryExpression}
 import org.apache.spark.sql.catalyst.expressions.Literal.TrueLiteral
 import org.apache.spark.sql.catalyst.plans.logical.{Assignment, Filter, 
Project, SupportsSubquery}
 import org.apache.spark.sql.execution.datasources.v2.DataSourceV2Relation
@@ -69,12 +69,17 @@ case class UpdatePaimonDataEvolutionTableCommand(
     val (updateTable, updateRelation) =
       MergeIntoPaimonDataEvolutionTable.withMatchedUpdateScanOptions(v2Table, 
relation)
     val targetRowId = rowIdAttribute(updateRelation)
-    val sourceTable = updatedRowIdSource(updateTable, updateRelation, 
targetRowId)
+    val matchedActionCondition = conditionAsMatchedAction
+    val sourceTable = updatedRowIdSource(
+      updateTable,
+      updateRelation,
+      targetRowId,
+      filterSource = !matchedActionCondition)
     val sourceRowId = sourceTable.output.head.asInstanceOf[AttributeReference]
 
     val matchedCondition = EqualTo(targetRowId, sourceRowId)
     val updateAction = SparkShimLoader.shim.createUpdateAction(
-      None,
+      if (matchedActionCondition) Some(condition) else None,
       alignedExpressions.map { case (expression, attribute) => 
Assignment(attribute, expression) })
 
     MergeIntoPaimonDataEvolutionTable(
@@ -124,20 +129,51 @@ case class UpdatePaimonDataEvolutionTableCommand(
     false
   }
 
+  /**
+   * Whether the WHERE condition is carried as the WHEN MATCHED condition of 
the self-merge instead
+   * of as a Filter on the source side.
+   *
+   * With the condition on the action, the source is `Project(PaimonRelation)` 
and
+   * [[MergeIntoPaimonDataEvolutionTable]] takes its self-merge shortcut: one 
scan of the files the
+   * condition can touch (pruned by file statistics), no self-join, no shuffle 
and no sort. Rows
+   * that fail the condition are copied through unchanged, which is exactly 
`UPDATE ... WHERE`.
+   *
+   * The Filter shape is kept when the condition cannot be evaluated inside 
`MergeRows`:
+   *   - a subquery can only be planned on a regular scan;
+   *   - attributes that do not belong to the relation (e.g. the read-side 
CHAR padding Project the
+   *     analyzer inserts on top of it) would be unresolved in the merge plan;
+   *   - a condition without column references (a constant, or e.g. `rand() < 
0.1`) gains nothing
+   *     from file pruning, and a constant-false one would rewrite every file 
as a no-op.
+   */
+  private def conditionAsMatchedAction: Boolean = {
+    if (condition == TrueLiteral || SubqueryExpression.hasSubquery(condition)) 
{
+      return false
+    }
+    val relationAttributes = relation.output ++ relation.metadataOutput
+    val references = condition.references.toSeq
+    references.nonEmpty &&
+    references.forall(attr => relationAttributes.exists(_.exprId == 
attr.exprId))
+  }
+
   private def updatedRowIdSource(
       updateTable: SparkTable,
       updateRelation: DataSourceV2Relation,
-      targetRowId: AttributeReference): Project = {
-    val conditionReferences = condition.references.toSeq.collect {
-      case attr: AttributeReference => attr
-    }
+      targetRowId: AttributeReference,
+      filterSource: Boolean): Project = {
+    val conditionReferences =
+      if (filterSource) {
+        condition.references.toSeq.collect { case attr: AttributeReference => 
attr }
+      } else {
+        Seq.empty
+      }
     val readOutput = deduplicateByExprId(conditionReferences :+ targetRowId)
     val sourceScan =
       SparkShimLoader.shim.copyDataSourceV2Relation(updateRelation, 
updateTable, readOutput)
-    // Keep the Filter visible for conditional UPDATEs. The data-evolution 
MERGE command uses a
-    // self-merge shortcut for Project(PaimonRelation); if a WHERE update were 
shaped that way, the
-    // shortcut would bypass the source join path and update every row.
-    val filteredSource = if (condition == TrueLiteral) sourceScan else 
Filter(condition, sourceScan)
+    // When the condition stays on the source side, keep the Filter visible: 
the data-evolution
+    // MERGE command uses a self-merge shortcut for Project(PaimonRelation), 
and that shortcut has
+    // no place for a source Filter, so it would update every row.
+    val filteredSource =
+      if (filterSource && condition != TrueLiteral) Filter(condition, 
sourceScan) else sourceScan
 
     Project(Seq(Alias(targetRowId, ROW_ID_COLUMN)()), filteredSource)
   }
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 a0a552baf4..0a5d30e231 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
@@ -1693,11 +1693,7 @@ abstract class RowTrackingTestBase extends 
PaimonSparkTestBase with AdaptiveSpar
 
         val (mergeRowsPlans, _) =
           executeMergeIntoAndCollectPlans("UPDATE t SET b = 22 WHERE id = 2")
-        assert(
-          mergeRowsPlans.exists(_.collectFirst { case _: Join => true 
}.nonEmpty),
-          s"Expected conditional UPDATE to use the general MERGE plan, but 
got: " +
-            mergeRowsPlans.mkString("\n")
-        )
+        assertSelfMergeShortcut(mergeRowsPlans)
         checkAnswer(
           sql("SELECT *, _ROW_ID, _SEQUENCE_NUMBER FROM t ORDER BY id"),
           Seq(Row(2, 22, 2, 0, 2), Row(3, 3, 3, 1, 2))
@@ -1706,6 +1702,153 @@ abstract class RowTrackingTestBase extends 
PaimonSparkTestBase with AdaptiveSpar
     }
   }
 
+  test("Data Evolution: V1 update with condition prunes files through the 
self-merge shortcut") {
+    withSparkSQLConf(
+      "spark.paimon.write.use-v2-write" -> "false",
+      "spark.paimon.data-evolution.merge-into.file-pruning" -> "true") {
+      withTable("target") {
+        sql("""
+              |CREATE TABLE target (id INT, b INT) TBLPROPERTIES (
+              |  'row-tracking.enabled' = 'true',
+              |  'data-evolution.enabled' = 'true',
+              |  'compaction.min.file-num' = '100')
+              |""".stripMargin)
+        sql("INSERT INTO target SELECT /*+ REPARTITION(1) */ * FROM VALUES (1, 
10), (2, 20)")
+        sql("INSERT INTO target SELECT /*+ REPARTITION(1) */ * FROM VALUES (3, 
30), (4, 40)")
+
+        // The condition reads the column being assigned: it must see the 
pre-update values.
+        val (mergeRowsPlans, resultedTableFiles) =
+          executeMergeIntoAndCollectPlans("UPDATE target SET b = b + 1 WHERE b 
>= 30")
+        assertSelfMergeShortcut(mergeRowsPlans)
+        assert(
+          resultedTableFiles == Seq(1L),
+          s"Expected the WHERE condition to prune the target scan to one file, 
but got: " +
+            resultedTableFiles.mkString(", "))
+        checkAnswer(
+          sql("SELECT id, b, _ROW_ID FROM target ORDER BY id"),
+          Seq(Row(1, 10, 0), Row(2, 20, 1), Row(3, 31, 2), Row(4, 41, 3)))
+
+        // A condition that matches nothing in the surviving file copies it 
through unchanged.
+        sql("UPDATE target SET b = 0 WHERE id = 4 AND b > 100")
+        checkAnswer(
+          sql("SELECT id, b FROM target ORDER BY id"),
+          Seq(Row(1, 10), Row(2, 20), Row(3, 31), Row(4, 41)))
+      }
+    }
+  }
+
+  test("Data Evolution: V1 update with partition condition prunes partitions 
through the shortcut") {
+    withSparkSQLConf(
+      "spark.paimon.write.use-v2-write" -> "false",
+      "spark.paimon.data-evolution.merge-into.file-pruning" -> "true") {
+      withTable("target") {
+        sql("""
+              |CREATE TABLE target (id INT, b INT, dt STRING) TBLPROPERTIES (
+              |  'row-tracking.enabled' = 'true',
+              |  'data-evolution.enabled' = 'true',
+              |  'compaction.min.file-num' = '100')
+              |PARTITIONED BY (dt)
+              |""".stripMargin)
+        sql("INSERT INTO target VALUES (1, 10, 'p1'), (2, 20, 'p1')")
+        sql("INSERT INTO target VALUES (3, 30, 'p2'), (4, 40, 'p2')")
+
+        val (mergeRowsPlans, resultedTableFiles) =
+          executeMergeIntoAndCollectPlans("UPDATE target SET b = 0 WHERE dt = 
'p2' AND id = 4")
+        assertSelfMergeShortcut(mergeRowsPlans)
+        assert(
+          resultedTableFiles == Seq(1L),
+          s"Expected the partition condition to prune the target scan to one 
file, but got: " +
+            resultedTableFiles.mkString(", "))
+        checkAnswer(
+          sql("SELECT id, b, dt FROM target ORDER BY id"),
+          Seq(Row(1, 10, "p1"), Row(2, 20, "p1"), Row(3, 30, "p2"), Row(4, 0, 
"p2")))
+      }
+    }
+  }
+
+  test("Data Evolution: V1 update condition on metadata column and NULL 
values") {
+    withSparkSQLConf("spark.paimon.write.use-v2-write" -> "false") {
+      withTable("t") {
+        sql(
+          "CREATE TABLE t (id INT, b INT) TBLPROPERTIES 
('row-tracking.enabled' = 'true', 'data-evolution.enabled' = 'true')")
+        sql("INSERT INTO t SELECT /*+ REPARTITION(1) */ * FROM VALUES (1, 10), 
(2, NULL), (3, 30)")
+
+        // A NULL condition value means "not matched": the row is copied 
through unchanged.
+        val (nullPlans, _) = executeMergeIntoAndCollectPlans("UPDATE t SET b = 
b + 1 WHERE b > 10")
+        assertSelfMergeShortcut(nullPlans)
+        checkAnswer(
+          sql("SELECT id, b FROM t ORDER BY id"),
+          Seq(Row(1, 10), Row(2, null), Row(3, 31)))
+
+        val (rowIdPlans, _) =
+          executeMergeIntoAndCollectPlans("UPDATE t SET b = 22 WHERE _ROW_ID = 
1")
+        assertSelfMergeShortcut(rowIdPlans)
+        checkAnswer(
+          sql("SELECT id, b, _ROW_ID FROM t ORDER BY id"),
+          Seq(Row(1, 10, 0), Row(2, 22, 1), Row(3, 31, 2)))
+      }
+    }
+  }
+
+  test("Data Evolution: V1 update with constant false condition changes 
nothing") {
+    withSparkSQLConf("spark.paimon.write.use-v2-write" -> "false") {
+      withTable("t") {
+        sql(
+          "CREATE TABLE t (id INT, b INT) TBLPROPERTIES 
('row-tracking.enabled' = 'true', 'data-evolution.enabled' = 'true')")
+        sql("INSERT INTO t VALUES (1, 10), (2, 20)")
+        val snapshotId = loadTable("t").snapshotManager().latestSnapshotId()
+
+        sql("UPDATE t SET b = 0 WHERE 1 = 0")
+
+        assert(loadTable("t").snapshotManager().latestSnapshotId() == 
snapshotId)
+        checkAnswer(sql("SELECT id, b FROM t ORDER BY id"), Seq(Row(1, 10), 
Row(2, 20)))
+      }
+    }
+  }
+
+  test("Data Evolution: V1 conditional update with user-specified snapshot 
uses the shortcut") {
+    withSparkSQLConf("spark.paimon.write.use-v2-write" -> "false") {
+      withTable("t") {
+        sql(
+          "CREATE TABLE t (id INT, b INT) TBLPROPERTIES 
('row-tracking.enabled' = 'true', 'data-evolution.enabled' = 'true')")
+        sql("INSERT INTO t VALUES (1, 10), (2, 20)")
+        val snapshotId = loadTable("t").snapshotManager().latestSnapshotId()
+        sql("INSERT INTO t VALUES (3, 30)")
+
+        var mergeRowsPlans = Seq.empty[LogicalPlan]
+        withSparkSQLConf("spark.paimon.scan.snapshot-id" -> 
snapshotId.toString) {
+          mergeRowsPlans = executeMergeIntoAndCollectPlans("UPDATE t SET b = 
100 WHERE id = 1")._1
+        }
+        assertSelfMergeShortcut(mergeRowsPlans)
+
+        checkAnswer(
+          sql("SELECT id, b FROM t ORDER BY id"),
+          Seq(Row(1, 100), Row(2, 20), Row(3, 30)))
+      }
+    }
+  }
+
+  test("Data Evolution: V1 update with subquery condition keeps the source 
filter") {
+    withSparkSQLConf("spark.paimon.write.use-v2-write" -> "false") {
+      withTable("t", "s") {
+        sql(
+          "CREATE TABLE t (id INT, b INT) TBLPROPERTIES 
('row-tracking.enabled' = 'true', 'data-evolution.enabled' = 'true')")
+        sql("INSERT INTO t VALUES (1, 10), (2, 20), (3, 30)")
+        sql("CREATE TABLE s (id INT)")
+        sql("INSERT INTO s VALUES (2), (3)")
+
+        val (mergeRowsPlans, _) =
+          executeMergeIntoAndCollectPlans("UPDATE t SET b = 0 WHERE id IN 
(SELECT id FROM s)")
+        assert(
+          mergeRowsPlans.exists(_.collectFirst { case _: Join => true 
}.nonEmpty),
+          s"Expected an UPDATE with a subquery condition to use the general 
MERGE plan, but got: " +
+            mergeRowsPlans.mkString("\n")
+        )
+        checkAnswer(sql("SELECT id, b FROM t ORDER BY id"), Seq(Row(1, 10), 
Row(2, 0), Row(3, 0)))
+      }
+    }
+  }
+
   test("Data Evolution: V1 update with global index updates unindexed rows") {
     withSparkSQLConf("spark.paimon.write.use-v2-write" -> "false") {
       withTable("t") {

Reply via email to