Repository navigation
feat: Push down offset/skip to TableScan - #25404
Conversation
Pass `offset` in addition to limit/fetch
|
Thank you for opening this pull request! Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch). Details |
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25404 +/- ##
==========================================
- Coverage 82.61% 82.58% -0.03%
==========================================
Files 1144 1144
Lines 442914 443198 +284
Branches 442914 443198 +284
==========================================
+ Hits 365898 366014 +116
- Misses 54862 54983 +121
- Partials 22154 22201 +47 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
Box TableScan's statistics_requests field because otherwise it becomes bigger than 176 bytes and `test_size_of_logical_plan()` fails
| /// Boxed to keep this rarely-populated field from growing every | ||
| /// `TableScan` (and thus `LogicalPlan`) by its own size; see | ||
| /// `test_size_of_logical_plan`. | ||
| pub statistics_requests: Box<BTreeSet<StatisticsRequest>>, |
There was a problem hiding this comment.
size_of::<BTreeSet<u32>>() => 24
size_of::<Box<BTreeSet<u32>>>() => 8
There was a problem hiding this comment.
verified that this is covered by an upgrade note
| .with_projection(projection) | ||
| .with_filters(filters) | ||
| .with_fetch(fetch) | ||
| .with_offset(offset) |
There was a problem hiding this comment.
It might be useful to add somewhere in the documentation what the expected order of offset/fetch wrt filters is. I assume it's the SQL order of filtering then offset/limit, but I think it would still be useful to make this explicit. Maybe a short note on the fields of TableScan?
There was a problem hiding this comment.
https://docs.rs/datafusion/latest/datafusion/catalog/trait.TableProvider.html#evaluation-order provides this information. I will update it with the skip field!
| /// Optional number of rows to read | ||
| pub fetch: Option<usize>, | ||
| /// Optional number of rows to skip | ||
| pub offset: Option<usize>, |
There was a problem hiding this comment.
Should we use 'skip' rather than 'offset' to be consistent with the rest of the code?
…ections (apache#25338) ## Which issue does this PR close? - Closes apache#25341. This PR supersedes apache#21363 by @crm26. It reuses the single mark join approach from that PR, and @crm26 is a co-author of the main commit. It also builds on apache#24972, which added the decorrelation of `IN` subqueries in a projection. ## Rationale for this change An `IN` or `NOT IN` subquery in a SELECT list gives correct results today, but the plan is quadratic. The optimizer builds three mark joins for each subquery. Two of them have no join predicate, so they run as nested loop joins over all outer rows and all inner rows. `COALESCE` duplicates the expression, so that shape gets six mark joins. This PR keeps one hash mark join for each subquery. The join is null-aware when a key can be NULL. A non-equality correlation (Q6 below) stays a residual filter of that join, which the null-aware hash join applies since apache#25339. | Shape | Query | main | this PR | Result rows identical | | --- | --- | --- | --- | --- | | Q1 | Bare `IN` in the SELECT list | 12.893 s | 0.005 s | yes | | Q2 | `COALESCE((x IN (...))::boolean, false)` | 36.430 s | 0.014 s | yes | | Q3 | Correlated `IN` with an equality predicate | 0.278 s | 0.008 s | yes | | Q4 | Two `IN` subqueries in separate columns | 20.904 s | 0.011 s | yes | | Q5 | Correlated `EXISTS` | 0.005 s | 0.003 s | yes | | Q6 | Correlated `IN` with a non-equality predicate | 2.489 s | 0.018 s | yes | The table is one run of the script below after the rebase onto `main` at `39ca2c74a9`, `datafusion-cli` built with `--profile release-nonlto` and the default `target_partitions`. Q6 in 3 more interleaved runs: `main` 2.57 s to 3.18 s, this PR 0.022 s to 0.058 s. These shapes are the `projection_subquery` benchmark suite from apache#25346. <details> <summary>Reproduction with datafusion-cli: script, plans and timings for each shape</summary> Run the script with `datafusion-cli -f mre.sql`. It creates two tables with 200000 rows each, then for each shape it prints `EXPLAIN` and runs the query. `datafusion-cli` prints the elapsed time after each statement. The result rows of every query are identical between the two builds. The plans and times for Q1 to Q5 below are from the first version of this PR, with `main` at 22651d2 and `target_partitions = 4`. The plans of this PR for those shapes did not change with the rebase. Q6 was captured again after the rebase. ```sql SET datafusion.execution.target_partitions = 4; CREATE TABLE outer_t AS SELECT CAST(v AS INT) AS id, CAST(v % 1000 AS INT) AS z FROM (SELECT unnest(generate_series(1, 200000)) AS v); CREATE TABLE inner_t AS SELECT CASE WHEN v % 97 = 0 THEN NULL ELSE CAST(v * 2 AS INT) END AS id, CAST(v % 1000 AS INT) AS z FROM (SELECT unnest(generate_series(1, 200000)) AS v); SELECT 'Q1' AS shape; EXPLAIN SELECT id, id IN (SELECT id FROM inner_t) AS m FROM outer_t; SELECT count(*) FILTER (WHERE m), count(*) FILTER (WHERE m IS NULL) FROM (SELECT id, id IN (SELECT id FROM inner_t) AS m FROM outer_t); SELECT 'Q2' AS shape; EXPLAIN SELECT id, COALESCE((id IN (SELECT id FROM inner_t))::boolean, false) AS m FROM outer_t; SELECT count(*) FILTER (WHERE m) FROM (SELECT id, COALESCE((id IN (SELECT id FROM inner_t))::boolean, false) AS m FROM outer_t); SELECT 'Q3' AS shape; EXPLAIN SELECT id, id IN (SELECT i.id FROM inner_t i WHERE i.z = o.z) AS m FROM outer_t o; SELECT count(*) FILTER (WHERE m), count(*) FILTER (WHERE m IS NULL) FROM (SELECT id, id IN (SELECT i.id FROM inner_t i WHERE i.z = o.z) AS m FROM outer_t o); SELECT 'Q4' AS shape; EXPLAIN SELECT id IN (SELECT id FROM inner_t) AS a, id IN (SELECT id FROM inner_t WHERE z < 500) AS b FROM outer_t; SELECT count(*) FILTER (WHERE a), count(*) FILTER (WHERE b) FROM (SELECT id IN (SELECT id FROM inner_t) AS a, id IN (SELECT id FROM inner_t WHERE z < 500) AS b FROM outer_t); SELECT 'Q5' AS shape; EXPLAIN SELECT id, EXISTS (SELECT 1 FROM inner_t i WHERE i.id = o.id) AS e FROM outer_t o; SELECT count(*) FILTER (WHERE e) FROM (SELECT id, EXISTS (SELECT 1 FROM inner_t i WHERE i.id = o.id) AS e FROM outer_t o); SELECT 'Q6' AS shape; EXPLAIN SELECT id, id IN (SELECT i.id FROM inner_t i WHERE i.z < o.z) AS m FROM outer_t o; SELECT count(*) FILTER (WHERE m), count(*) FILTER (WHERE m IS NULL) FROM (SELECT id, id IN (SELECT i.id FROM inner_t i WHERE i.z < o.z) AS m FROM outer_t o); ``` <details> <summary>Q1: Bare `IN` in the SELECT list, plans on main and on this PR</summary> **main**, query time 190.730 s ``` logical_plan Projection: outer_t.id, __correlated_sq_1.mark IS NOT DISTINCT FROM Boolean(true) OR (__correlated_sq_2.mark OR outer_t.id IS NULL AND __correlated_sq_3.mark) IS NOT DISTINCT FROM Boolean(true) AND __correlated_sq_1.mark IS DISTINCT FROM Boolean(true) AND Boolean(NULL) AS m LeftMark Join: LeftMark Join: LeftMark Join: outer_t.id = __correlated_sq_1.id null_aware TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_2 Filter: inner_t.id IS NULL TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_3 TableScan: inner_t projection=[id] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 IS NOT DISTINCT FROM true OR (mark@2 OR id@0 IS NULL AND mark@3) IS NOT DISTINCT FROM true AND mark@1 IS DISTINCT FROM true AND NULL as m] NestedLoopJoinExec: join_type=LeftMark CoalescePartitionsExec NestedLoopJoinExec: join_type=RightMark CoalescePartitionsExec FilterExec: id@0 IS NULL DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` **this PR**, query time 0.034 s ``` logical_plan Projection: outer_t.id, __correlated_sq_1.mark AS m LeftMark Join: outer_t.id = __correlated_sq_1.id null_aware TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 TableScan: inner_t projection=[id] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 as m] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` </details> <details> <summary>Q2: `COALESCE((x IN (...))::boolean, false)`, plans on main and on this PR</summary> **main**, query time 351.579 s ``` logical_plan Projection: outer_t.id, CAST(__correlated_sq_1.mark IS NOT DISTINCT FROM Boolean(true) OR (__correlated_sq_2.mark OR outer_t.id IS NULL AND __correlated_sq_3.mark) IS NOT DISTINCT FROM Boolean(true) AND __correlated_sq_1.mark IS DISTINCT FROM Boolean(true) AND Boolean(NULL) AS Boolean) IS NOT NULL AND CAST(__correlated_sq_4.mark IS NOT DISTINCT FROM Boolean(true) OR (__correlated_sq_5.mark OR outer_t.id IS NULL AND __correlated_sq_6.mark) IS NOT DISTINCT FROM Boolean(true) AND __correlated_sq_4.mark IS DISTINCT FROM Boolean(true) AND Boolean(NULL) AS Boolean) AS m LeftMark Join: LeftMark Join: LeftMark Join: outer_t.id = __correlated_sq_4.id null_aware LeftMark Join: LeftMark Join: LeftMark Join: outer_t.id = __correlated_sq_1.id null_aware TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_2 Filter: inner_t.id IS NULL TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_3 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_4 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_5 Filter: inner_t.id IS NULL TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_6 TableScan: inner_t projection=[id] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 IS NOT DISTINCT FROM true OR (mark@2 OR id@0 IS NULL AND mark@3) IS NOT DISTINCT FROM true AND mark@1 IS DISTINCT FROM true AND NULL IS NOT NULL AND (mark@4 IS NOT DISTINCT FROM true OR (mark@5 OR id@0 IS NULL AND mark@6) IS NOT DISTINCT FROM true AND mark@4 IS DISTINCT FROM true AND NULL) as m] NestedLoopJoinExec: join_type=LeftMark CoalescePartitionsExec NestedLoopJoinExec: join_type=RightMark CoalescePartitionsExec FilterExec: id@0 IS NULL DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec NestedLoopJoinExec: join_type=LeftMark CoalescePartitionsExec NestedLoopJoinExec: join_type=RightMark CoalescePartitionsExec FilterExec: id@0 IS NULL DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` **this PR**, query time 0.054 s ``` logical_plan Projection: outer_t.id, CAST(__correlated_sq_1.mark AS Boolean) IS NOT NULL AND CAST(__correlated_sq_2.mark AS Boolean) AS m LeftMark Join: outer_t.id = __correlated_sq_2.id null_aware LeftMark Join: outer_t.id = __correlated_sq_1.id null_aware TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_2 TableScan: inner_t projection=[id] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 IS NOT NULL AND mark@2 as m] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` </details> <details> <summary>Q3: Correlated `IN` with an equality predicate, plans on main and on this PR</summary> **main**, query time 1.372 s ``` logical_plan Projection: o.id, __correlated_sq_1.mark IS NOT DISTINCT FROM Boolean(true) OR (__correlated_sq_2.mark OR o.id IS NULL AND __correlated_sq_3.mark) IS NOT DISTINCT FROM Boolean(true) AND __correlated_sq_1.mark IS DISTINCT FROM Boolean(true) AND Boolean(NULL) AS m LeftMark Join: o.z = __correlated_sq_3.z LeftMark Join: o.z = __correlated_sq_2.z LeftMark Join: o.id = __correlated_sq_1.id, o.z = __correlated_sq_1.z null_aware SubqueryAlias: o TableScan: outer_t projection=[id, z] SubqueryAlias: __correlated_sq_1 SubqueryAlias: i TableScan: inner_t projection=[id, z] SubqueryAlias: __correlated_sq_2 SubqueryAlias: i Projection: inner_t.z Filter: inner_t.id IS NULL TableScan: inner_t projection=[id, z] SubqueryAlias: __correlated_sq_3 SubqueryAlias: i TableScan: inner_t projection=[z] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 IS NOT DISTINCT FROM true OR (mark@2 OR id@0 IS NULL AND mark@3) IS NOT DISTINCT FROM true AND mark@1 IS DISTINCT FROM true AND NULL as m] HashJoinExec: mode=CollectLeft, join_type=RightMark, on=[(z@0, z@1)], projection=[id@0, mark@2, mark@3, mark@4] CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] HashJoinExec: mode=CollectLeft, join_type=RightMark, on=[(z@0, z@1)] CoalescePartitionsExec FilterExec: id@0 IS NULL, projection=[z@1] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0), (z@1, z@1)], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` **this PR**, query time 0.091 s ``` logical_plan Projection: o.id, __correlated_sq_1.mark AS m LeftMark Join: o.id = __correlated_sq_1.id, o.z = __correlated_sq_1.z null_aware SubqueryAlias: o TableScan: outer_t projection=[id, z] SubqueryAlias: __correlated_sq_1 SubqueryAlias: i TableScan: inner_t projection=[id, z] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 as m] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0), (z@1, z@1)], projection=[id@0, mark@2], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` </details> <details> <summary>Q4: Two `IN` subqueries in separate columns, plans on main and on this PR</summary> **main**, query time 296.676 s ``` logical_plan Projection: __correlated_sq_1.mark IS NOT DISTINCT FROM Boolean(true) OR (__correlated_sq_2.mark OR outer_t.id IS NULL AND __correlated_sq_3.mark) IS NOT DISTINCT FROM Boolean(true) AND __correlated_sq_1.mark IS DISTINCT FROM Boolean(true) AND Boolean(NULL) AS a, __correlated_sq_4.mark IS NOT DISTINCT FROM Boolean(true) OR (__correlated_sq_5.mark OR outer_t.id IS NULL AND __correlated_sq_6.mark) IS NOT DISTINCT FROM Boolean(true) AND __correlated_sq_4.mark IS DISTINCT FROM Boolean(true) AND Boolean(NULL) AS b LeftMark Join: LeftMark Join: LeftMark Join: outer_t.id = __correlated_sq_4.id null_aware LeftMark Join: LeftMark Join: LeftMark Join: outer_t.id = __correlated_sq_1.id null_aware TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_2 Filter: inner_t.id IS NULL TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_3 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_4 Projection: inner_t.id Filter: inner_t.z < Int32(500) TableScan: inner_t projection=[id, z] SubqueryAlias: __correlated_sq_5 Projection: inner_t.id Filter: inner_t.z < Int32(500) AND inner_t.id IS NULL TableScan: inner_t projection=[id, z] SubqueryAlias: __correlated_sq_6 Projection: inner_t.id Filter: inner_t.z < Int32(500) TableScan: inner_t projection=[id, z] physical_plan ProjectionExec: expr=[mark@1 IS NOT DISTINCT FROM true OR (mark@2 OR id@0 IS NULL AND mark@3) IS NOT DISTINCT FROM true AND mark@1 IS DISTINCT FROM true AND NULL as a, mark@4 IS NOT DISTINCT FROM true OR (mark@5 OR id@0 IS NULL AND mark@6) IS NOT DISTINCT FROM true AND mark@4 IS DISTINCT FROM true AND NULL as b] NestedLoopJoinExec: join_type=LeftMark CoalescePartitionsExec NestedLoopJoinExec: join_type=RightMark CoalescePartitionsExec FilterExec: z@1 < 500 AND id@0 IS NULL, projection=[id@0] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec NestedLoopJoinExec: join_type=LeftMark CoalescePartitionsExec NestedLoopJoinExec: join_type=RightMark CoalescePartitionsExec FilterExec: id@0 IS NULL DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] FilterExec: z@1 < 500, projection=[id@0] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] FilterExec: z@1 < 500, projection=[id@0] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` **this PR**, query time 0.063 s ``` logical_plan Projection: __correlated_sq_1.mark AS a, __correlated_sq_2.mark AS b LeftMark Join: outer_t.id = __correlated_sq_2.id null_aware LeftMark Join: outer_t.id = __correlated_sq_1.id null_aware TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 TableScan: inner_t projection=[id] SubqueryAlias: __correlated_sq_2 Projection: inner_t.id Filter: inner_t.z < Int32(500) TableScan: inner_t projection=[id, z] physical_plan ProjectionExec: expr=[mark@0 as a, mark@1 as b] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], projection=[mark@1, mark@2], null_aware CoalescePartitionsExec HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], null_aware CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] FilterExec: z@1 < 500, projection=[id@0] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` </details> <details> <summary>Q5: Correlated `EXISTS`, plans on main and on this PR</summary> **main**, query time 0.021 s ``` logical_plan Projection: o.id, __correlated_sq_1.mark AS e LeftMark Join: o.id = __correlated_sq_1.id SubqueryAlias: o TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 SubqueryAlias: i TableScan: inner_t projection=[id] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 as e] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)] CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` **this PR**, query time 0.021 s ``` logical_plan Projection: o.id, __correlated_sq_1.mark AS e LeftMark Join: o.id = __correlated_sq_1.id SubqueryAlias: o TableScan: outer_t projection=[id] SubqueryAlias: __correlated_sq_1 SubqueryAlias: i TableScan: inner_t projection=[id] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 as e] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)] CoalescePartitionsExec DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] DataSourceExec: partitions=4, partition_sizes=[7, 6, 6, 6] ``` </details> <details> <summary>Q6: Correlated `IN` with a non-equality predicate, plans on main and on this PR</summary> **main** (`39ca2c74a9`), query time 2.489 s ``` logical_plan Projection: o.id, CASE WHEN __correlated_sq_1.mark THEN Boolean(true) WHEN __correlated_sq_2.mark OR o.id IS NULL AND __correlated_sq_3.mark THEN Boolean(NULL) ELSE Boolean(false) END AS m LeftMark Join: Filter: __correlated_sq_3.z < o.z LeftMark Join: Filter: __correlated_sq_2.z < o.z LeftMark Join: o.id = __correlated_sq_1.id Filter: __correlated_sq_1.z < o.z null_aware SubqueryAlias: o TableScan: outer_t projection=[id, z] SubqueryAlias: __correlated_sq_1 SubqueryAlias: i TableScan: inner_t projection=[id, z] SubqueryAlias: __correlated_sq_2 SubqueryAlias: i Projection: inner_t.z Filter: inner_t.id IS NULL TableScan: inner_t projection=[id, z] SubqueryAlias: __correlated_sq_3 SubqueryAlias: i TableScan: inner_t projection=[z] physical_plan ProjectionExec: expr=[id@0 as id, CASE WHEN mark@1 THEN true WHEN mark@2 OR id@0 IS NULL AND mark@3 THEN NULL ELSE false END as m] NestedLoopJoinExec: join_type=LeftMark, filter=z@1 < z@0, projection=[id@0, mark@2, mark@3, mark@4] CoalescePartitionsExec NestedLoopJoinExec: join_type=RightMark, filter=z@1 < z@0 CoalescePartitionsExec FilterExec: id@0 IS NULL, projection=[z@1] DataSourceExec: partitions=12, partition_sizes=[3, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], filter=z@1 < z@0, null_aware CoalescePartitionsExec DataSourceExec: partitions=12, partition_sizes=[3, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2] DataSourceExec: partitions=12, partition_sizes=[3, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2] DataSourceExec: partitions=12, partition_sizes=[3, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2] ``` **this PR**, query time 0.018 s ``` logical_plan Projection: o.id, __correlated_sq_1.mark AS m LeftMark Join: o.id = __correlated_sq_1.id Filter: __correlated_sq_1.z < o.z null_aware SubqueryAlias: o TableScan: outer_t projection=[id, z] SubqueryAlias: __correlated_sq_1 SubqueryAlias: i TableScan: inner_t projection=[id, z] physical_plan ProjectionExec: expr=[id@0 as id, mark@1 as m] HashJoinExec: mode=CollectLeft, join_type=LeftMark, on=[(id@0, id@0)], filter=z@1 < z@0, projection=[id@0, mark@2], null_aware CoalescePartitionsExec DataSourceExec: partitions=12, partition_sizes=[3, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2] DataSourceExec: partitions=12, partition_sizes=[3, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2] ``` </details> </details> ## What changes are included in this PR? 1. `build_join` now reports whether the mark column of the `LeftMark` join is exact under three-valued logic. It is exact when no side of the `IN` predicate can be NULL in scope, or when the join is null-aware and the `IN` equality is its first key. A residual non-equality filter does not change this, because the null-aware hash join applies it when it decides whether a NULL makes the mark UNKNOWN (apache#25339). `in_subquery_value_mark_join` builds this join first. If the mark is exact, it returns the mark column alone, or `NOT mark` for `NOT IN`. The three join materialization stays for the other cases. 2. The rule decides null-awareness from one question: can the `IN` value or the subquery output be NULL inside the scope of an outer row? Only then can `IN` be UNKNOWN. - The test uses the key expressions, not the columns they reference. A key such as `NULLIF(id, 1)` or `TRY_CAST(s AS INT)` can be NULL over a `NOT NULL` column. - `PullUpCorrelatedExpr` now records every correlated conjunct, before it drops the conjunct that repeats the `IN` predicate. A conjunct such as `y = x`, `y > x` or `y IS NOT NULL` is never TRUE for a NULL `y`, so it keeps a NULL `y` out of the scope. The conjuncts keep their outer references, so a subquery column with the same qualified name as an outer column is not mistaken for the outer value. - A correlation key that is NULL only empties the scope, so it never makes the join null-aware. - The recorded conjuncts are cleared when the pull up passes an outer join, a union or a grouping set, because such a node can put a NULL back into the column. With `ROLLUP`, the grand-total row is a NULL for every outer row. 3. A `NOT IN` filter uses the same scope-aware nullability for its null-aware `LeftAnti` join. If the `IN` equality does not become the first key of that join, the rule gives up and the `NOT IN` becomes the mark joins, which give the correct result. 4. The regression test for apache#24574 in `projection_pushdown.slt` now uses a correlated subquery with `LIMIT 1`. The subquery then still reaches `ExtractLeafExpressions`, and the alias generator still starts at 2 inside it. The earlier version of the test let the subquery be flattened, so it no longer exercised the scan inside a subquery path. 5. New sqllogictest cases in `subquery_projection.slt`: - `EXPLAIN` guards: one hash mark join for each subquery, one null-aware mark join with a residual filter, a null-aware `LeftAnti` join with two keys, and with one key and a residual filter. - NULL semantics: `IN` and `NOT IN` with a NULL in the subquery, a NULL outer value, `NOT` inside `CASE`, a projection over an aggregate, the `COALESCE` and cast shape, and the residual filter shape. - Nullable key expressions over `NOT NULL` columns: a `NULLIF` value, a `TRY_CAST` value and a `NULLIF` subquery output, in a projection and in a `WHERE` clause. - A correlation that repeats the `IN` predicate, and a subquery column with the same qualified name as the outer value. - The same correlation below a `ROLLUP`, as a projected `IN` with an `EXPLAIN` guard and as a `NOT IN` filter. The expected results agree with DuckDB 1.5.2, and the earlier cases also with PostgreSQL. 6. The `projection_subquery` benchmark docs no longer describe q07 as the nested-loop control. ## DuckDB comparison For an absolute reference, the same seven queries in `datafusion-cli` and in DuckDB 1.5.2, on the same two tables, release build, Apple M4 Pro. `main` is the merge-base `64871d9` and "this PR" is `3af87370c1`. The later commits of this PR change the correlated and `NOT IN` filter paths, and, after the rebase onto apache#25339, the plan of q07. q07 is measured again below. The tables are the ones the suite builds, at its default of 30,000 rows each. Each number is the median of 7 runs. One round runs one session per engine, and the rounds interleave the three engines. | Query | `main` | this PR | DuckDB | | --- | --- | --- | --- | | q01 bare `IN` | 533 ms | 2 ms | 2 ms | | q02 `COALESCE` over `IN` | 1352 ms | 3 ms | 2 ms | | q03 correlated `IN`, equality | 14 ms | 3 ms | 3 ms | | q04 two `IN` columns | 1026 ms | 3 ms | 2 ms | | q05 bare `NOT IN` | 646 ms | 1 ms | 2 ms | | q06 correlated `EXISTS` | 4 ms | 2 ms | 2 ms | | q07 correlated `IN`, non-equality | 117 ms | 93 ms | 458 ms | All three engines give the same result for all seven queries, and the counts are the ones in the suite's checked-in result files. For the four quadratic shapes, q01, q02, q04 and q05, `main` is 270x to 680x slower than DuckDB. This PR puts them at DuckDB's cost. These times are at the 1 ms resolution of both command line tools, so the exact factor is approximate; the size of the difference is not. q07 now uses one null-aware mark join with the `<` correlation as a residual filter. Measured again after the rebase, median of 7 runs at 30,000 rows: `main` at `39ca2c74a9` 78 ms, this PR 3 ms, DuckDB 404 ms. The counts are the same in all three. A second run at 100,000 rows per table gives the same picture: `main` takes 2741 ms to 6355 ms for q01, q02, q04 and q05, this PR takes 2 ms to 4 ms and DuckDB 3 ms to 4 ms. All three engines again give the same counts. <details><summary>How to run it</summary> The queries are the suite's own `benchmarks/sql_benchmarks/projection_subquery/queries/q0*.sql`, with an alias added to the derived table. The load SQL needs one change for DuckDB, because `generate_series` gives a column of that name there and a column named `value` in DataFusion: ```sql CREATE TABLE outer_t AS SELECT CAST(value AS INT) AS id, CAST(value % 1000 AS INT) AS z FROM generate_series(1, 30000) AS g(value); CREATE TABLE inner_t AS SELECT CASE WHEN value % 97 = 0 THEN NULL ELSE CAST(value * 2 AS INT) END AS id, CAST(value % 1000 AS INT) AS z FROM generate_series(1, 30000) AS g(value); ``` `datafusion-cli` reports the time of each statement as `Elapsed`, and `duckdb` reports it as `Run Time (s): real` after `.timer on`. </details> ## What is the testing strategy for this PR? - Unit tests in `datafusion/optimizer/src/decorrelate_predicate_subquery.rs`. The updated snapshots show a single mark join. New tests cover a residual filter, `NOT IN`, nullable key expressions, a correlation that repeats the `IN` predicate, and the same correlation below a `ROLLUP`. The last one fails without the clearing in item 2. - The new sqllogictest cases in `subquery_projection.slt` listed above. - The full sqllogictest suite, the optimizer crate and workspace clippy are green locally. - The benchmark numbers above. ## Are there any user-facing changes? Plans for projected `IN` and `NOT IN` subqueries are much faster. There are no changes to any public API, and no documentation change is needed. Some query results change, and each change is a correction. DuckDB 1.5.2 gives the new results: - An `IN` or `NOT IN` whose key is an expression that can be NULL over a `NOT NULL` column, such as `NULLIF(id, 1)`. For a projected `IN`, `main` returns `false` instead of `NULL`. For a `NOT IN` filter, `main` keeps rows that it must drop. The correlated form of that is apache#25347. - A correlated `NOT IN` whose correlation repeats the `IN` predicate, `x NOT IN (SELECT y FROM t WHERE y = x)`. `main` drops rows that it must keep (apache#25480). - A projected `IN` whose correlation repeats the `IN` predicate below a `ROLLUP`, `k IN (SELECT i.k FROM i WHERE i.k = o.k GROUP BY ROLLUP(i.k))`. `main` fails to plan it since apache#25529. It now gives NULL for a miss, as DuckDB does. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
## Which issue does this PR close? - Part of apache#25035. Split out from apache#25771 as requested, so the evaluator change there can be benchmarked against `main`. ## Rationale for this change On `main`, a left-deep `AND` chain can be several times slower than the same chain built right-deep, once the prefix drops below the pre-selection threshold (apache#25035). The existing `binary_op` benchmarks don't cover this. ## What changes are included in this PR? A `conjunction` group in `datafusion/physical-expr/benches/binary_op.rs`: left-deep and right-deep `AND` chains of cheap comparisons, casts, `IN` lists and arithmetic, with and without unused columns, plus nullable and regex-suffix variants. Only this file changes; no engine code. ## What is the testing strategy for this PR? Benchmark only. All `binary_op` benchmarks run on `main` (`cargo test --bench binary_op`). Run with `cargo bench --bench binary_op -- conjunction`. ## Are there any user-facing changes? No.
…n SLT (apache#25651) ## Which issue does this PR close? - Part of apache#25650. ## Rationale for this change `MemoryPool` accounting only covers memory that operators explicitly reserve, so a process can be OOM killed while the pool reports plenty of headroom. As a first step, this makes the gap visible: it logs the difference between what pools have reserved and what's actually allocated, and it's on by default in SLT. It only logs and never fails a test, unlike apache#22626 which was reverted in apache#22860. ## What changes are included in this PR? - `datafusion-execution`: `MemoryDriftTracker` and `DriftLoggingPool`. The pool wraps any `MemoryPool` and reports reservation changes to a tracker, which can be shared by many pools. The tracker compares the reserved total against allocated bytes from a caller-supplied `Fn() -> usize`, since DataFusion doesn't choose the global allocator. It logs at `info` each time positive drift rises by 64 MB, naming the pool and consumer, and records the peak. - `sqllogictest`: a `CountingAllocator` global allocator, with every test file's pool wrapped against one process-wide tracker. The peak drift is printed at the end of the run. Disable it with `--memory-drift false`. Test files run concurrently, so the comparison is process-wide rather than per file. Files that `SET datafusion.runtime.memory_limit` replace their pool and drop out of the reserved total (6 of 504 files). Example output from a full local run: ``` Peak memory drift: drift=84.3 MB allocated=85.4 MB reserved=1176.3 KB pool=order.slt consumer=ExternalSorterMerge[0] ``` The allocator batches its counts per thread, so SLT runtime is unchanged locally (7s both with and without it). ## Are these changes tested? Yes. There are unit tests for the tracker and pool, plus a doc test. I also ran the full SLT suite locally with drift logging on and off. ## Are there any user-facing changes? New public types in `datafusion_execution::memory_pool`. No changes to existing APIs.
…25854) ## Which issue does this PR close? No linked issue. This is a follow-up to the fully matched Parquet row-group work in PR apache#23696. ## Rationale for this change When row-group statistics prove that every row satisfies a scan predicate, a Bloom filter cannot prune that group. Reading its Bloom filters still adds object-store I/O, which can be especially costly for remote files. ## What changes are included in this PR? - Skip Bloom filter reads and predicate evaluation for fully matched row groups. - Avoid creating a Bloom reader when every surviving row group is fully matched. - Keep the Bloom pruning matched metric accounting for skipped groups. ## What is the testing strategy for this PR? The new `fully_matched_row_groups_skip_bloom_filter_reads` test verifies that a partially matched group still reads Bloom filters, a fully matched group reduces `bytes_scanned`, and an all-fully-matched file reads zero Bloom bytes during open. It also checks that the returned rows are unchanged. ## Are there any user-facing changes? No API or query-result changes. Scans avoid unnecessary Bloom filter reads for fully matched row groups.
## Which issue does this PR close? Related to apache#25620. Split from apache#25648 following [the review request](apache#25648 (comment)). This PR can be reviewed and merged independently. ## Rationale for this change A pushed-down `fetch` stops each FilterExec partition after a bounded number of rows, but filter statistics currently ignore it. A filter with 100 matching rows and `fetch=3` still estimates 100 rows in a single partition. This can distort downstream planning decisions. ## What changes are included in this PR? Apply the per-partition fetch in both the built-in filter statistics and FilterStatisticsProvider. Adjust row and byte estimates while keeping empty, all-null, and proven singleton column statistics consistent. Overall statistics remain estimates when rows may be distributed unevenly across partitions. Add independent execution tests with one and two partitions, plus tests for zero fetch, all-null columns, singleton preservation, and the provider path. ## What is the testing strategy for this PR? Passed on this independent branch: - `cargo test --profile ci -p datafusion-physical-plan --lib` (2,367 tests) - `cargo clippy --profile ci -p datafusion-physical-plan --all-targets --all-features -- -D warnings` - `cargo fmt --all` and `git diff --check` Benchmarks were not rerun for the split. This extracts the fetch implementation already present in apache#25648; its earlier combined measurements are not isolated measurements of this PR. ## Are there any user-facing changes? EXPLAIN statistics account for a fetch pushed into FilterExec. Plan choices can change as a result. No public API changes.
It is too slow and it did not reveal any problems Suggested-by @adriangb at apache#25404 (review)
Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com>
Which issue does this PR close?
Rationale for this change
Currently
TableProviderprovides thelimitargument to thescan()method to request a maximum number of rows. There is no way to tell the implementation to skip some of the rows, for example to fetch the second/third/Nth page of rows (i.e. SQL... LIMIT 20 OFFSET 40).Adding an additional field to ScanArgs (named
offsetorskip) will make it possible for implementations to override thescan_with_args()method and optimize their scan to read and return only the requested rows.What changes are included in this PR?
offsetis added toScanArgs, with a setter and a getter.TableProvidertrait -supports_offset_pushdown() -> bool. By default it returnsfalsebut any implementation that can support skipping of rows could override it to returntrueand combined with a custom implementation ofscan_with_args()to optimise its data scan/read.TableProvider::scan()to use::scan_with_args()where they could support offset push downNote:
datafusion-ffiis not updated because it does not exposescan_with_args()yet.What is the testing strategy for this PR?
New unit tests are added for the implementations which support offset pushdown.
Are there any user-facing changes?
The new functionality is opt-in! All currently existing custom implementations of
TableProvidertrait will continue to compile and run without any modifications.Any custom implementation that wants to make use of the new functionality will need to override
TableProvider::supports_offset_pushdown()to returntrueand make use ofScanArgs::offsetin itsscan_with_args()implementation.