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