This is an automated email from the ASF dual-hosted git repository.
YannByron 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 e0e7ac9de0 [spark] Fix MERGE INTO silently skipping WHEN NOT MATCHED
BY SOURCE actions (#9155)
e0e7ac9de0 is described below
commit e0e7ac9de023cce6609b7f10ae1288e8b7c3401c
Author: Kerwin Zhang <[email protected]>
AuthorDate: Tue Aug 11 14:52:22 2026 +0800
[spark] Fix MERGE INTO silently skipping WHEN NOT MATCHED BY SOURCE actions
(#9155)
---
.../spark/commands/MergeIntoPaimonTable.scala | 16 +++++++++-
.../spark/commands/MergeIntoPaimonTable.scala | 16 +++++++++-
.../sql/MergeIntoNotMatchedBySourceTest.scala | 37 ++++++++++++++++++++++
3 files changed, 67 insertions(+), 2 deletions(-)
diff --git
a/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
b/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
index c24d4d6675..d51504059d 100644
---
a/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
+++
b/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
@@ -60,8 +60,22 @@ case class MergeIntoPaimonTable(
lazy val relation: DataSourceV2Relation =
PaimonRelation.getPaimonRelation(targetTable)
+ /**
+ * The target-only part of the merge condition, used to prune the target
table before the join.
+ *
+ * Pruning is only sound when the merge has no `WHEN NOT MATCHED BY SOURCE`
action. A target row
+ * that fails the target-only condition can never satisfy the whole merge
condition, so it can
+ * never be matched — dropping it therefore cannot change the outcome of the
`WHEN MATCHED` and
+ * `WHEN NOT MATCHED` actions. Those very rows are, however, exactly the
population that
+ * `WHEN NOT MATCHED BY SOURCE` is defined over, so pruning them would
silently skip the actions
+ * that should apply to them. Disable the pruning in that case.
+ */
private lazy val (targetOnlyCondition, filteredTargetPlan):
(Option[Expression], LogicalPlan) = {
- val filtersOnlyTarget = getExpressionOnlyRelated(mergeCondition,
targetTable)
+ val filtersOnlyTarget = if (notMatchedBySourceActions.isEmpty) {
+ getExpressionOnlyRelated(mergeCondition, targetTable)
+ } else {
+ None
+ }
(
filtersOnlyTarget,
filtersOnlyTarget
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
index c24d4d6675..d51504059d 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonTable.scala
@@ -60,8 +60,22 @@ case class MergeIntoPaimonTable(
lazy val relation: DataSourceV2Relation =
PaimonRelation.getPaimonRelation(targetTable)
+ /**
+ * The target-only part of the merge condition, used to prune the target
table before the join.
+ *
+ * Pruning is only sound when the merge has no `WHEN NOT MATCHED BY SOURCE`
action. A target row
+ * that fails the target-only condition can never satisfy the whole merge
condition, so it can
+ * never be matched — dropping it therefore cannot change the outcome of the
`WHEN MATCHED` and
+ * `WHEN NOT MATCHED` actions. Those very rows are, however, exactly the
population that
+ * `WHEN NOT MATCHED BY SOURCE` is defined over, so pruning them would
silently skip the actions
+ * that should apply to them. Disable the pruning in that case.
+ */
private lazy val (targetOnlyCondition, filteredTargetPlan):
(Option[Expression], LogicalPlan) = {
- val filtersOnlyTarget = getExpressionOnlyRelated(mergeCondition,
targetTable)
+ val filtersOnlyTarget = if (notMatchedBySourceActions.isEmpty) {
+ getExpressionOnlyRelated(mergeCondition, targetTable)
+ } else {
+ None
+ }
(
filtersOnlyTarget,
filtersOnlyTarget
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoNotMatchedBySourceTest.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoNotMatchedBySourceTest.scala
index d2c822789c..0ab6e5755a 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoNotMatchedBySourceTest.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/MergeIntoNotMatchedBySourceTest.scala
@@ -220,4 +220,41 @@ trait MergeIntoNotMatchedBySourceTest extends
PaimonSparkTestBase with PaimonTab
}
}
}
+
+ test("Paimon MergeInto: not matched by source is not narrowed by target-only
condition") {
+ withTable("source", "target") {
+
+ Seq((1, 100)).toDF("a", "b").createOrReplaceTempView("source")
+
+ createTable("target", "a INT, b INT, c STRING, pt STRING", Seq("a",
"pt"), Seq("pt"))
+ spark.sql("""
+ |INSERT INTO target VALUES
+ | (1, 10, 'c1', 'p1'), (2, 20, 'c2', 'p1'),
+ | (3, 30, 'c3', 'p2'), (4, 40, 'c4', 'p2')
+ |""".stripMargin)
+
+ // `t.pt = 'p1'` only references the target, so it is a candidate for
pruning the target
+ // before the join. Pruning it away would also drop the 'p2' rows from
the population that
+ // WHEN NOT MATCHED BY SOURCE is defined over, silently skipping their
update.
+ spark.sql("""
+ |MERGE INTO target t
+ |USING source s
+ |ON t.a = s.a AND t.pt = 'p1'
+ |WHEN MATCHED THEN
+ | UPDATE SET t.b = s.b
+ |WHEN NOT MATCHED BY SOURCE THEN
+ | UPDATE SET t.c = 'stale'
+ |""".stripMargin)
+
+ // a=1: matched (pt='p1') => b updated to 100
+ // a=2: not matched by source (pt='p1') => c = 'stale'
+ // a=3, a=4: not matched by source (pt='p2', excluded by the target-only
condition, so no
+ // source row can ever match them) => c = 'stale'
+ checkAnswer(
+ spark.sql("SELECT a, b, c, pt FROM target ORDER BY a"),
+ Row(1, 100, "c1", "p1") :: Row(2, 20, "stale", "p1") ::
+ Row(3, 30, "stale", "p2") :: Row(4, 40, "stale", "p2") :: Nil
+ )
+ }
+ }
}