Repository navigation
Conversation
…runtime filtering ### What changes were proposed in this pull request? Scalar subquery filters on partition columns (e.g., `WHERE d_date_sk = (SELECT min(d_date_sk) FROM ...)`) are excluded from pushdown in DSv2 at every stage. The filter lands as a `FilterExec` above `BatchScanExec`, evaluated row-by-row. The scan reads all partitions -- no partition pruning occurs. DSv1 already handles this: `FileSourceStrategy` puts subquery filters in `partitionFilters`, `isDynamicFilter` classifies them as dynamic, and `getPartitionPruningFilterFromBroadcast` calls `ScalarSubquery.toLiteral` at execution time for partition pruning via `listFiles()`. This PR routes partition-column scalar subquery filters into `BatchScanExec.runtimeFilters`, leveraging the existing `SupportsRuntimeV2Filtering.filter()` infrastructure: - **DataSourceV2Strategy**: When the scan implements `SupportsRuntimeV2Filtering`, extract subquery filters from `postScanFilters` where references are a subset of partition columns. Add to `runtimeFilters` alongside existing DPP filters. They remain in `postScanFilters` as a correctness safety net (V2 `filter()` is advisory). - **BatchScanExec**: In `filteredPartitions`, non-DPP runtime filters are literalized (replacing `ExecScalarSubquery` with its resolved literal) and translated to V2 predicates via `translateFilterV2`. - **InMemoryTableWithV2Filter** (test infra): Added `=` predicate handling in `filter()` alongside existing `IN`, plus a `case _ =>` catch-all. No new interfaces, no config flags, no connector changes needed. ### Why are the changes needed? TPC-DS queries with scalar subquery partition filters (e.g., Q5, Q12, Q16, Q20, Q37, Q77, Q80, Q92, Q94, Q95) read all partitions in DSv2 scans even though the subquery resolves to a single value at runtime. This causes significant I/O overhead that DSv1 avoids. ### Does this PR introduce _any_ user-facing change? No API changes. Queries with scalar subquery filters on partition columns will now benefit from partition pruning in DSv2 scans, reducing I/O. ### How was this patch tested? New unit test in `DataSourceV2SQLSuiteV2Filter`: - Creates a 10-partition table and a dimension table - Runs `SELECT * FROM t WHERE part = (SELECT max(val) FROM dim)` - Asserts query correctness, scalar subquery presence in `runtimeFilters`, and exactly 1 partition after pruning ### Was this patch authored or co-authored using generative AI tooling? Yes, co-authored with Claude Code.
cloud-fan
left a comment
There was a problem hiding this comment.
Review Summary
Prior state and problem: DSv2 scans (BatchScanExec) only supported partition pruning via Dynamic Partition Pruning (DPP). Scalar subquery filters on partition columns (e.g., WHERE d_date_sk = (SELECT min(d_date_sk) FROM ...)) were evaluated row-by-row by FilterExec above the scan, with no partition pruning. DSv1 already handled this — FileSourceStrategy includes scalar subquery partition filters in partitionFilters, evaluated at runtime for file listing.
Design approach: Route scalar subquery partition filters into the existing runtimeFilters infrastructure. This mirrors the DSv1 approach while respecting the V2 architecture (where partition pruning goes through the connector's SupportsRuntimeV2Filtering.filter() API). The filters remain in postScanFilters as a safety net — the V2 filter() is advisory, so FilterExec still evaluates them. Subquery reuse (ReuseExchangeAndSubquery / ReuseAdaptiveSubquery) ensures the duplicated subquery expression is only executed once.
Implementation sketch:
DataSourceV2Strategy: Extracts scalar subquery filters whose column references are a subset offilterAttributes(partition columns). Appends them toruntimeFiltersalongside DPP filters.BatchScanExec.filteredPartitions: New path for non-DPP runtime filters — literalizesExecScalarSubquery→Literal, then translates to V2PredicateviatranslateFilterV2.InMemoryTableWithV2Filter(test infra): Adds=predicate handling alongside existingIN, plus acase _ =>catch-all.- Test: Verifies query correctness, scalar subquery presence in
runtimeFilters, and partition count after pruning.
- Use containsPattern(SCALAR_SUBQUERY) instead of hasSubquery to match only scalar subqueries, consistent with DSv1 path in FileSourceStrategy - Add NULL guard in InMemoryTableWithV2Filter equality predicate to avoid NPE on null literal values - Replace placeholder SPARK-XXXXX with SPARK-56467 in test name Co-authored-by: Isaac
|
Thanks for the review! Addressed all three points:
|
Co-authored-by: Isaac
…nslateScalarSubqueryFilterV2 utility Co-authored-by: Isaac
|
@szehon-ho addressed comments, thanks! |
Co-authored-by: Isaac
|
thanks, merging to master! |
What changes were proposed in this pull request?
Scalar subquery filters on partition columns (e.g.,
WHERE d_date_sk = (SELECT min(d_date_sk) FROM ...)) are excluded from pushdown in DSv2 at every stage. The filter lands as aFilterExecaboveBatchScanExec, evaluated row-by-row. The scan reads all partitions -- no partition pruning occurs.DSv1 already handles this:
FileSourceStrategyputs subquery filters inpartitionFilters,isDynamicFilterclassifies them as dynamic, andgetPartitionPruningFilterFromBroadcastcallsScalarSubquery.toLiteralat execution time for partition pruning vialistFiles().This PR routes partition-column scalar subquery filters into
BatchScanExec.runtimeFilters, leveraging the existingSupportsRuntimeV2Filtering.filter()infrastructure:SupportsRuntimeV2Filtering, extract subquery filters frompostScanFilterswhere references are a subset of partition columns. Add toruntimeFiltersalongside existing DPP filters. They remain inpostScanFiltersas a correctness safety net (V2filter()is advisory).filteredPartitions, non-DPP runtime filters are literalized (replacingExecScalarSubquerywith its resolved literal) and translated to V2 predicates viatranslateFilterV2.=predicate handling infilter()alongside existingIN, plus acase _ =>catch-all.No new interfaces, no config flags, no connector changes needed.
Why are the changes needed?
TPC-DS queries with scalar subquery partition filters (e.g., Q5, Q12, Q16, Q20, Q37, Q77, Q80, Q92, Q94, Q95) read all partitions in DSv2 scans even though the subquery resolves to a single value at runtime. This causes significant I/O overhead that DSv1 avoids.
Does this PR introduce any user-facing change?
No API changes. Queries with scalar subquery filters on partition columns will now benefit from partition pruning in DSv2 scans, reducing I/O.
How was this patch tested?
New unit test in
DataSourceV2SQLSuiteV2Filter:SELECT * FROM t WHERE part = (SELECT max(val) FROM dim)runtimeFilters, and exactly 1 partition after pruningWas this patch authored or co-authored using generative AI tooling?
Yes, co-authored with Claude Code.