From e09b3adfcb8508584984b0e11d0d9a9575a0d375 Mon Sep 17 00:00:00 2001 From: Peter Toth Date: Tue, 4 Aug 2026 12:45:31 +0200 Subject: [PATCH 1/2] [SPARK-58207][SQL][FOLLOWUP] Skip runtime filter pushdown for nondeterministic filters DataSourceV2Strategy no longer routes a non-deterministic post-scan filter into BatchScanExec.runtimeFilters, so it is not pushed to a SupportsRuntimeV2Filtering scan at runtime. #57357 guarded only PushDownUtils.pushFilters (pushdown at query compilation); the runtime path had no determinism guard, and since a scalar subquery runtime filter also stays in postScanFilters, the predicate was always evaluated twice with different results. --- docs/sql-migration-guide.md | 2 +- .../datasources/v2/DataSourceV2Strategy.scala | 7 +++- .../sql/connector/DataSourceV2SQLSuite.scala | 41 +++++++++++++++++++ 3 files changed, 48 insertions(+), 2 deletions(-) diff --git a/docs/sql-migration-guide.md b/docs/sql-migration-guide.md index f25295ad5349e..67ba504ad6f02 100644 --- a/docs/sql-migration-guide.md +++ b/docs/sql-migration-guide.md @@ -34,7 +34,7 @@ license: | - Since Spark 4.3, `spark.sql.execution.replaceHashWithSortAgg` defaults to `true`. Spark now replaces a hash-based aggregate with a sort aggregate when the aggregate's child is already sorted on the grouping keys. To restore the previous behavior, set `spark.sql.execution.replaceHashWithSortAgg` to `false`. - Since Spark 4.3, `spark.sql.execution.combineAdjacentAggregation` defaults to `true`. Spark now merges an adjacent partial/final aggregate pair (with no shuffle between them) into a single complete-mode aggregate. This setting is independent of `spark.sql.execution.replaceHashWithSortAgg`, so disabling only `replaceHashWithSortAgg` still leaves adjacent aggregation combined; to fully restore the previous partial/final staging, set both `spark.sql.execution.replaceHashWithSortAgg` and `spark.sql.execution.combineAdjacentAggregation` to `false`. - Since Spark 4.3, the exact `percentile`, `percentile_cont`, and `median` aggregate functions (including their `WITHIN GROUP (ORDER BY ...)` forms) compute the linear interpolation between two neighboring values as `lower + fraction * (higher - lower)` instead of `(1 - fraction) * lower + fraction * higher`. The two are equal in exact arithmetic, but the new form is monotonically non-decreasing in the requested percentage and avoids a rounding error the old form could introduce. As a result these functions may return a value that differs from earlier releases in the last ULP. `percentile_disc` and `percentile_approx` are unaffected. -- Since Spark 4.3, non-deterministic filters (for example predicates involving `rand()`) are no longer pushed down to DataSource V2 sources that implement `SupportsPushDownV2Filters`; they are evaluated by Spark after the scan instead. This prevents a source from evaluating such a predicate a different number of times than Spark, or using it for pruning while also returning it for post-scan re-evaluation. +- Since Spark 4.3, non-deterministic filters (for example predicates involving `rand()`) are no longer pushed down to DataSource V2 sources: neither at query compilation to sources that implement `SupportsPushDownV2Filters`, nor as runtime filters to sources that implement `SupportsRuntimeV2Filtering`; they are evaluated by Spark after the scan instead. This prevents a source from evaluating such a predicate a different number of times than Spark, or using it for pruning while also returning it for post-scan re-evaluation. - Since Spark 4.3, the new `COMMENT ON COLUMN ... IS NULL` syntax removes a column comment by passing a `null` comment to `TableChange.updateColumnComment(String[], String)`, so a `UpdateColumnComment` table change may now carry a `null` `newComment()`. Previously `newComment()` was always non-null. DataSource V2 catalogs that handle `UpdateColumnComment` should null-check `newComment()` and treat `null` as "remove the column comment" (mirroring how `UpdateColumnDefaultValue` already carries a `null` value to drop a default). - Since Spark 4.3, `unix_seconds`, `unix_millis`, and `unix_micros` accept `TIMESTAMP_NTZ` and the nanosecond-precision timestamp types directly, reading them with no time-zone shift. Previously these functions accepted only `TIMESTAMP_LTZ`; a `TIMESTAMP_NTZ` or nanosecond-timestamp argument was rejected with a `DATATYPE_MISMATCH` error. This is a new capability and does not change the result of any query that previously succeeded. - Since Spark 4.3, `hash()` and `xxhash64()` include the `days` field of `CalendarInterval` when computing the hash, so their output for interval values differs from earlier releases. Previously the codegen path dropped `days`, disagreeing with interpreted evaluation. diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala index 72807a5242c7c..21e230f61e5aa 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala @@ -172,9 +172,14 @@ class DataSourceV2Strategy(session: SparkSession) extends Strategy with Predicat // These filters stay in postScanFilters for correctness (FilterExec above scan), // but are also routed into runtimeFilters so BatchScanExec can use them for // partition pruning via SupportsRuntimeV2Filtering.filter(). + // Non-deterministic filters are not routed: they would be pushed to the source for + // pruning while the FilterExec above the scan re-evaluates them, so the two evaluations + // would disagree and rows the source pruned away could not be recovered. This is the + // runtime counterpart of the pushFilters guard in PushDownUtils (SPARK-58207). val scalarSubqueryFilters = if (relation.runtimeFilterAttrs.nonEmpty) { postScanFilters.filter { f => - f.containsPattern(SCALAR_SUBQUERY) && + f.deterministic && + f.containsPattern(SCALAR_SUBQUERY) && f.references.nonEmpty && f.references.subsetOf(relation.runtimeFilterAttrs) } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/DataSourceV2SQLSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/DataSourceV2SQLSuite.scala index 79c46bdb4ccf7..76c7902188ed3 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/DataSourceV2SQLSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/DataSourceV2SQLSuite.scala @@ -5478,6 +5478,47 @@ class DataSourceV2SQLSuiteV2Filter extends DataSourceV2SQLSuite { s"Expected 1 partition after scalar subquery pruning, got $numPartitions") } } + + test("SPARK-58207: non-deterministic scalar subquery filters are not pushed into " + + "runtimeFilters") { + val tbl = s"${catalogAndNamespace}tbl" + val dim = s"${catalogAndNamespace}dim" + withTable(tbl, dim) { + sql(s"CREATE TABLE $tbl (id INT, part INT) USING $v2Format PARTITIONED BY (part)") + for (i <- 0 until 10) { + sql(s"INSERT INTO $tbl VALUES ($i, $i)") + } + + sql(s"CREATE TABLE $dim (val INT) USING $v2Format") + sql(s"INSERT INTO $dim VALUES (3)") + + // `part = (subquery) OR rand() < 0.5` references only the partition column and holds a + // scalar subquery, so it is a candidate for runtime pushdown, but it is non-deterministic. + // Routing it would push it to the source for pruning while the FilterExec above the scan + // re-evaluates it, and a partition the source dropped on its own evaluation could not be + // recovered. + val df = sql( + s"SELECT * FROM $tbl WHERE part = (SELECT max(val) FROM $dim) OR rand() < 0.5") + df.collect() + + val batchScan = collect(df.queryExecution.executedPlan) { + case b: BatchScanExec => b + }.head + assert(batchScan.runtimeFilters.isEmpty, + s"Expected no runtime filters for a non-deterministic filter, " + + s"got ${batchScan.runtimeFilters}") + + // No pruning at the source, and the filter is still evaluated by Spark after the scan. + val numPartitions = batchScan.filteredPartitions.count(_.isDefined) + assert(numPartitions == 10, + s"Expected all 10 partitions to be retained, got $numPartitions") + val postScanConditions = collect(df.queryExecution.executedPlan) { + case f: FilterExec => f.condition + } + assert(postScanConditions.exists(!_.deterministic), + s"Expected the non-deterministic filter above the scan, got $postScanConditions") + } + } } class ReserveSchemaNullabilityCatalog extends InMemoryCatalog { From 96d0fee0508d626ccb7413f2420764164511b2f2 Mon Sep 17 00:00:00 2001 From: Peter Toth Date: Tue, 4 Aug 2026 13:31:35 +0200 Subject: [PATCH 2/2] Address review: DPP bypass rationale, caller contract, wording - Cite CleanupDynamicPruningFilters as the reason dynamicFilters needs no determinism check (a non-deterministic DynamicPruningSubquery is already rewritten to TrueLiteral before planning), rather than the weaker "not duplicated" argument. - Document the caller contract on PushDownUtils.pushRuntimeFilters, since the invariant is now enforced two files away in DataSourceV2Strategy. - Soften "would disagree" to "may disagree": the source is free to ignore a pushed runtime filter. --- .../sql/execution/datasources/v2/DataSourceV2Strategy.scala | 5 ++++- .../spark/sql/execution/datasources/v2/PushDownUtils.scala | 5 +++++ 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala index 21e230f61e5aa..ce9092ea810b5 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2Strategy.scala @@ -174,7 +174,7 @@ class DataSourceV2Strategy(session: SparkSession) extends Strategy with Predicat // partition pruning via SupportsRuntimeV2Filtering.filter(). // Non-deterministic filters are not routed: they would be pushed to the source for // pruning while the FilterExec above the scan re-evaluates them, so the two evaluations - // would disagree and rows the source pruned away could not be recovered. This is the + // may disagree and rows the source pruned away could not be recovered. This is the // runtime counterpart of the pushFilters guard in PushDownUtils (SPARK-58207). val scalarSubqueryFilters = if (relation.runtimeFilterAttrs.nonEmpty) { postScanFilters.filter { f => @@ -186,6 +186,9 @@ class DataSourceV2Strategy(session: SparkSession) extends Strategy with Predicat } else { Seq.empty } + // dynamicFilters need no such check: a DynamicPruningSubquery over a non-deterministic + // filtering plan is itself non-deterministic, so CleanupDynamicPruningFilters has already + // rewritten it to TrueLiteral by the time we get here. val runtimeFilters = dynamicFilters ++ scalarSubqueryFilters val batchExec = BatchScanExec(relation.output, relation.scan, runtimeFilters, diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/PushDownUtils.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/PushDownUtils.scala index 46c4e98595827..5e1e22a0c66f0 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/PushDownUtils.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/PushDownUtils.scala @@ -170,6 +170,11 @@ object PushDownUtils extends Logging { * Note: Do not call multiple times for the same `scan` instance; * [[SupportsRuntimeV2Filtering.filter]] is mutating. * + * Note: `runtimeFilters` must not contain non-deterministic filters. A runtime filter is also + * evaluated by the `FilterExec` above the scan, so pushing a non-deterministic one would + * evaluate it twice with different results. `DataSourceV2Strategy` enforces this where + * `runtimeFilters` is built (SPARK-58207). + * * @return true if any filters were pushed to the data source */ def pushRuntimeFilters(