Repository navigation
perf: use one null-aware mark join for hashable IN subqueries in projections - #25338
Conversation
9e94b59 to
04f3d7d
Compare
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #25338 +/- ##
==========================================
- Coverage 82.57% 82.57% -0.01%
==========================================
Files 1142 1142
Lines 440958 441219 +261
Branches 440958 441219 +261
==========================================
+ Hits 364141 364322 +181
- Misses 54806 54854 +48
- Partials 22011 22043 +32 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@kosiew since you reviewed #24972 would you mind taking a look here? cc @neilconway since you were involved in #21363 |
|
run benchmarks |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark clickbench_partitionedResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark tpcdsResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark tpchResults will be posted here when complete File an issue against this benchmark runner |
There was a problem hiding this comment.
🟡 Changes recommended
Expression-level nullability can be missed, causing projected IN to return false instead of UNKNOWN.
Get a fresh assessment by requesting another Copilot review.
Pull request overview
Optimizes projected IN/NOT IN subqueries by using a single null-aware mark join when possible.
Changes:
- Tracks whether mark joins preserve three-valued logic.
- Retains three-join fallback for residual predicates.
- Adds plan and NULL-semantics regression tests.
File summaries
| File | Description |
|---|---|
datafusion/optimizer/src/decorrelate_predicate_subquery.rs |
Implements single-mark-join optimization. |
datafusion/sqllogictest/test_files/subquery_projection.slt |
Tests plans and NULL semantics. |
datafusion/sqllogictest/test_files/projection_pushdown.slt |
Updates alias-collision regression coverage. |
Review details
- Files reviewed: 3/3 changed files
- Comments generated: 1
- Review effort level: Balanced
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark tpchCPU Details (lscpu)Details
Resource Usagetpch — base (merge-base)
tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark tpcdsCPU Details (lscpu)Details
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark clickbench_partitionedCPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
run benchmark clickbench_partitioned |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing 22651d2 (22651d2) to 22651d2 (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
ref: "22651d24cc8196f3206e09a36a437eda9e766b87"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing 22651d2 (22651d2) to 22651d2 (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
ref: "22651d24cc8196f3206e09a36a437eda9e766b87"CPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
run benchmark clickbench_partitioned |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark clickbench_partitionedResults will be posted here when complete File an issue against this benchmark runner |
|
I opened #25346 to add benchmarks |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing in-exists-subquery-projection (04f3d7d) to 22651d2 (merge-base) diff Run configurationrun benchmark clickbench_partitionedCPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
kosiew
left a comment
There was a problem hiding this comment.
Thanks for working on this. Reusing the exact hashable LeftMark join for projected IN / NOT IN is a nice improvement, and the added NULL-semantics coverage is helpful.
I found one blocking issue in the correlated NOT IN path. The expression-level nullability change can now make a LeftAnti join null-aware when there is more than one hash key, but physical planning only supports a single key for null-aware LeftAnti joins. I left a repro and suggested adding an execution regression test below.
I also left one non-blocking suggestion for the empty-subquery boundary on the new single-mark-join path.
| && join_keys_may_be_null(&join_filter, left.schema(), sub_query_alias.schema())?; | ||
| // Additionally, if no join key can be NULL on either side, we don't need | ||
| // null-aware semantics because NULLs cannot exist in the keys. | ||
| let null_aware = if join_type == JoinType::LeftAnti && in_predicate_opt.is_some() { |
There was a problem hiding this comment.
I think this introduces a planning failure for correlated NOT IN when the nullable expression key is combined with another correlated equality key.
For example, SELECT id FROM o WHERE NULLIF(id, 1) NOT IN (SELECT id FROM r WHERE r.grp = o.grp) with non-nullable o(id, grp) and r(id, grp) now makes the LeftAnti join null-aware because NULLIF(id, 1) is nullable. The correlation adds grp as a second hash key, so physical planning rejects the resulting join with null_aware LeftAnti joins only support single column join key, got 2 columns.
This looks newly reachable through the expression-level nullability change. Could we preserve the correct correlated NOT IN semantics using a supported fallback or plan shape here? It would also be good to add this query as an execution regression test and assert that only the SQL-true rows are returned.
There was a problem hiding this comment.
Thanks, confirmed. The query failed with got 2 columns. It plans again in fe408e3: a LeftAnti join with more than one key keeps the column test from main.
This does not give the correct result for your NULLIF example. The join has two keys and is not null-aware, so the k = 1 row stays, the same as on main. The test in subquery_projection.slt records this, with a comment that links apache/datafusion#25347. A correct plan needs a null-aware LeftAnti join with more than one key. apache/datafusion#25339 (approved) adds that. When it merges, I will remove this special case and update the expected output to the correct rows. I did not want to add a second fallback plan here, because apache/datafusion#25339 is the fix for this.
| 05)----ProjectionExec: expr=[CAST(id@0 AS Int64) as r3.id] | ||
| 06)------DataSourceExec: partitions=1, partition_sizes=[1] | ||
|
|
||
| # `NULLIF(id, 1)` is NULL for `id = 1`, and `r3` has no NULL, so the answer is |
There was a problem hiding this comment.
Could we also add an empty-subquery case for the new single null-aware mark path using a typed nullable key expression, for example NULLIF(id, 1)?
In particular, it would be useful to assert that the NULL-key row produces false, not NULL, when the subquery is empty, and that the plan still uses a single mark join. The existing top-level NULL IN (empty) test covers the SQL result semantics, but it takes the legacy three-join path, so it does not protect this boundary of the new optimization.
|
Thanks @adriangb , here is a suggestion: Correlated The For a correlated CREATE TABLE t1(k INT NOT NULL, s VARCHAR NOT NULL) AS VALUES (1, 'a');
CREATE TABLE t2(k INT NOT NULL, s VARCHAR NOT NULL) AS VALUES (1, 'B');
SELECT * FROM t1 WHERE upper(t1.s) NOT IN (SELECT t2.s FROM t2 WHERE t2.k = t1.k);
-- main: plans a plain LeftAnti and returns (1, 'a')
-- this PR: expected "null_aware LeftAnti joins only support single column join key"Nullable columns already fail like this on main (#25347), but this change extends the failure to common function keys over non-nullable data. It also pins uncorrelated Suggest limiting the expression-level check on the let null_aware = if join_type == JoinType::LeftAnti && in_predicate_opt.is_some() {
let (equijoin_keys, residual_filter) = split_eq_and_noneq_join_predicate(
join_filter.clone(),
left.schema(),
sub_query_alias.schema(),
)?;
- join_keys_may_be_null(
- &equijoin_keys,
- residual_filter.as_ref(),
- left.schema(),
- sub_query_alias.schema(),
- )?
+ if equijoin_keys.len() == 1 {
+ join_keys_may_be_null(
+ &equijoin_keys,
+ residual_filter.as_ref(),
+ left.schema(),
+ sub_query_alias.schema(),
+ )?
+ } else {
+ // Null-aware LeftAnti supports a single key only (#25347); keep the
+ // previous column-based test so correlated NOT IN still plans.
+ join_filter_columns_may_be_null(&join_filter, left.schema(), sub_query_alias.schema())?
+ }
} else {Please also add an slt case with the correlated query above. |
|
Thanks @jayzhan211, good catch. I applied your suggestion in fe408e3. |
|
@jayzhan211 @kosiew could we merge the benchmarks in #25346 before this change so we can look at perf numbers? |
984bfbd to
fb08efa
Compare
jayzhan211
left a comment
There was a problem hiding this comment.
@adriangb I leave some suggestions
| .all(|&e| can_pullup_over_aggregation(e)); | ||
| let (mut join_filters, subquery_filters) = | ||
| find_join_exprs(subquery_filter_exprs)?; | ||
| for expr in &join_filters { |
There was a problem hiding this comment.
correlated_filters comes from find_join_exprs, which strips outer refs → a subquery column sharing the outer value's qualified name is read as the value → join not null-aware → wrong result (regression vs main).
Repro (expected NULL for the NULL row; PR gives false, main gives NULL; the NOT IN filter keeps the NULL row):
create table o(x int, k int) as values (null, 1), (7, 1), (8, 1);
create table t(x int, y int not null) as values (1, 7);
select a.x, a.x in (select a.y from t as a where a.x = b.k) as r
from o as a, (select 1 as k) as b order by 1;
select a.x from o as a, (select 1 as k) as b
where a.x not in (select a.y from t as a where a.x = b.k);Fix (tested: repro correct, optimizer tests + subquery_projection/null_aware*/joins/subquery slt pass): keep the outer refs and match per side.
- let (mut join_filters, subquery_filters) =
- find_join_exprs(subquery_filter_exprs)?;
- for expr in &join_filters {
- if !self.correlated_filters.contains(expr) {
- self.correlated_filters.push(expr.clone());
+ for expr in &subquery_filter_exprs {
+ if expr.contains_outer() && !self.correlated_filters.contains(expr) {
+ self.correlated_filters.push((*expr).clone());
}
}
+ let (mut join_filters, subquery_filters) =
+ find_join_exprs(subquery_filter_exprs)?;fn filter_rejects_null(filter: &Expr, key: &Expr, key_is_outer: bool) -> bool {
let is_key = |side: &Expr| {
let side = strip_casts(side);
if key_is_outer {
side.contains_outer()
&& side.column_refs().is_empty()
&& &strip_outer_reference(side.clone()) == key
} else {
!side.contains_outer() && side == key
}
};
match filter {
Expr::BinaryExpr(BinaryExpr { left, op, right }) => {
matches!(
op,
Operator::Eq
| Operator::NotEq
| Operator::Lt
| Operator::LtEq
| Operator::Gt
| Operator::GtEq
) && (is_key(left) || is_key(right))
}
Expr::IsNotNull(expr) => is_key(expr),
_ => false,
}
}Pass true for value_as_written and false for output_expr from InValue::may_be_null_in_scope, and add the repro to subquery_projection.slt.
There was a problem hiding this comment.
Thank you, this is a real regression. I took your fix in 10ccc3e. correlated_filters now keeps the outer references, and filter_rejects_null matches each side of a conjunct against the side of the join that the key is on. Your two queries are in subquery_projection.slt (tables sh_o and sh_t), and DuckDB 1.5.2 gives the same results: NULL for the NULL row, and 8 for the NOT IN filter.
| // row", which is the weaker fact that this case needs, and which the first | ||
| // reading implies. `IS NULL` is two-valued on both sides, so neither test | ||
| // adds an UNKNOWN of its own. | ||
| let unknown_alias = alias.next("__correlated_sq"); |
There was a problem hiding this comment.
Commit 50d6e4a3a (one merged UNKNOWN join) slows the residual fallback ~1.6x: the NLJ now evaluates k < k AND (y IS NULL OR x IS NULL) on every pair; before, one NLJ ran bare k < k and the other had its right side pre-filtered to y IS NULL.
release-nonlto, 200k x 200k, 3 interleaved runs: main 2.80/3.55/3.78 s · bb043f10e (parent) 3.05/3.06/3.67 s · PR 4.74/5.14/6.48 s. EXPLAIN ANALYZE NLJ elapsed_compute: main 21.7 s + 2.7 s, PR 44.7 s.
create table outer_t as select i as id, case when i % 10 = 0 then null else i end as x, i % 100 as k from generate_series(1, 200000) t(i);
create table inner_t as select i as id, case when i % 7 = 0 then null else i * 2 end as y, i % 100 as k from generate_series(1, 200000) t(i);
select id, x in (select y from inner_t where inner_t.k < outer_t.k) as r from outer_t;Your Q6 number went the other way, so this is data-dependent. Please drop the commit from this PR (the fallback then matches main) and revisit it with #25336.
There was a problem hiding this comment.
Agreed. I removed 50d6e4a from this PR. I also removed 1a9d842, which only fixed an error message that the merged join caused. The fallback for a residual filter is now the same plan as on main.
After the rebase onto #25339, a NOT IN filter with a residual filter no longer uses the fallback. It becomes one null-aware anti join, because that join now applies the residual filter. The nai_res_og plan that #25339 pins stays as it is, and the joins.slt plan is the same as on main again.
1a9d842 to
5148cd3
Compare
| # `NULLIF(k, 1)` is NULL for `k = 1`, and that group of `t2` is not empty, so | ||
| # the correct result has no row for `k = 1`. The join has two keys, the value | ||
| # and the correlation, and the null-aware `LeftAnti` executor takes one key | ||
| # only. The `NOT IN` becomes a null-aware mark join, which takes any number of | ||
| # keys, and a filter on the mark. `main` keeps the `k = 1` row, which is | ||
| # https://github.com/apache/datafusion/issues/25347. | ||
| query I rowsort | ||
| SELECT k FROM t1 WHERE NULLIF(t1.k, 1) NOT IN (SELECT t2.k + 10 FROM t2 WHERE t2.k = t1.k); | ||
| ---- | ||
| 2 |
There was a problem hiding this comment.
Result changed here, and the old value was wrong. This expectation was 1 and 2; it is now 2 only. DuckDB 1.5.2 agrees with the new value.
NULLIF(t1.k, 1) is NULL for the row k = 1. The correlated subquery for that row gives {11}, which is not empty, so NULL NOT IN ({11}) is UNKNOWN and the row must not appear.
The old plan could not say that. A null-aware LeftAnti join takes one key only, and this join has two, the value and the correlation, so the join stayed a plain anti join and kept the row. The NOT IN now becomes a null-aware LeftMark join, which takes any number of keys, plus a filter on the mark.
D SELECT t1.k, NULLIF(t1.k,1) AS key, (SELECT list(t2.k+10) FROM t2 WHERE t2.k=t1.k) AS sub FROM t1 ORDER BY t1.k;
┌───────┬───────┬───────────┐
│ k │ key │ sub │
├───────┼───────┼───────────┤
│ 1 │ NULL │ [11] │
│ 2 │ 2 │ [12] │
└───────┴───────┴───────────┘
D SELECT k FROM t1 WHERE NULLIF(t1.k, 1) NOT IN (SELECT t2.k + 10 FROM t2 WHERE t2.k = t1.k);
┌───┐
│ k │
├───┤
│ 2 │
└───┘
This closes the part of #25347 that a single key cannot express.
| # The same shape where the correlated subquery result is not empty for the NULL | ||
| # key: `NULL NOT IN ({5})` is UNKNOWN, so `2` is the only row. `main` builds a | ||
| # plain anti join here and also keeps `1`. | ||
| query I rowsort | ||
| SELECT k FROM ra WHERE NULLIF(ra.k, 1) NOT IN (SELECT rb.k FROM rb WHERE rb.z > ra.z); | ||
| ---- | ||
| 2 |
There was a problem hiding this comment.
Result changed here, and the old value was wrong. This expectation was 1 and 2; it is now 2 only. DuckDB 1.5.2 agrees with the new value.
NULLIF(ra.k, 1) is NULL for the row k = 1, and the correlated subquery for that row gives {5}, which is not empty. NULL NOT IN ({5}) is UNKNOWN, so the row must not appear.
A residual filter stays on this join, and no hash join can mark the UNKNOWN rows of a residual filter. The NOT IN now becomes the three mark joins that materialize its three-valued result, the plan that #24972 already uses for a projected IN.
D SELECT ra.k, NULLIF(ra.k,1) AS key, (SELECT list(rb.k) FROM rb WHERE rb.z > ra.z) AS sub FROM ra ORDER BY ra.k;
┌───────┬───────┬──────────┐
│ k │ key │ sub │
├───────┼───────┼──────────┤
│ 1 │ NULL │ [5] │
│ 2 │ 2 │ [5] │
└───────┴───────┴──────────┘
D SELECT k FROM ra WHERE NULLIF(ra.k, 1) NOT IN (SELECT rb.k FROM rb WHERE rb.z > ra.z);
┌───┐
│ k │
├───┤
│ 2 │
└───┘
The row above, with rb.z < ra.z, keeps its expectation of 1 and 2. Its subquery result is empty for the NULL key, and NULL NOT IN (<empty set>) is TRUE.
This is the executor gap in #25336. #25339 fixes the executor, and this shape can go back to one anti join when that lands.
| # The correlation names no column of the subquery, so it stays as a residual | ||
| # filter and the subquery is either the whole of `ic` or empty. A NULL in | ||
| # `ic.id` then makes the answer UNKNOWN only for the rows whose correlation | ||
| # holds. `main` builds a null-aware anti join that does not apply the residual | ||
| # when it looks for a NULL, sees one in `ic`, and drops every row. This is the | ||
| # shape whose plan `joins.slt` pins. | ||
| statement ok | ||
| CREATE TABLE oc(id INT, g INT) AS VALUES (1, 1), (2, 0), (NULL, 0); | ||
|
|
||
| statement ok | ||
| CREATE TABLE ic(id INT) AS VALUES (5), (NULL); | ||
|
|
||
| # `g > 0` holds for `id = 1` only, so its subquery is `{5, NULL}` and | ||
| # `1 NOT IN {5, NULL}` is UNKNOWN. The other two rows have an empty subquery, | ||
| # and `<anything> NOT IN (<empty set>)` is TRUE, the NULL row included. | ||
| query I rowsort | ||
| SELECT id FROM oc WHERE oc.id NOT IN (SELECT ic.id FROM ic WHERE oc.g > 0); | ||
| ---- | ||
| 2 | ||
| NULL |
There was a problem hiding this comment.
New test, and it is the one that shows what the joins.slt plan change buys. main returns no rows for this query. DuckDB 1.5.2 returns the two rows below.
The correlation oc.g > 0 names no column of ic, so it cannot become a join key and stays as a residual filter. The subquery is therefore the whole of ic for an outer row whose g > 0, and empty for every other row.
D SELECT oc.id, oc.g, (SELECT list(ic.id) FROM ic WHERE oc.g > 0) AS sub FROM oc ORDER BY oc.id;
┌───────┬───────┬─────────────┐
│ id │ g │ sub │
├───────┼───────┼─────────────┤
│ 1 │ 1 │ [5, NULL] │
│ 2 │ 0 │ NULL │
│ NULL │ 0 │ NULL │
└───────┴───────┴─────────────┘
D SELECT id FROM oc WHERE oc.id NOT IN (SELECT ic.id FROM ic WHERE oc.g > 0);
┌───────┐
│ id │
├───────┤
│ 2 │
│ NULL │
└───────┘
main builds one null-aware anti join and keeps the residual filter on it. That executor does not apply the residual when it looks for a NULL, so it finds the NULL in ic and treats every outer row as UNKNOWN, including the two whose subquery is empty. It returns nothing.
The NULL row is kept on purpose: its subquery is empty, and NULL NOT IN (<empty set>) is TRUE.
This is the same shape as the joins.slt query below, which had no result test.
| # A grouping set above the correlated filter is a different problem. The | ||
| # grand-total row of `ROLLUP` exists for every outer row, so a miss is UNKNOWN | ||
| # here, not FALSE. The pull up moves the filter above the aggregate, which | ||
| # loses that row; the same plan gives a wrong `EXISTS` too. That is a bug in | ||
| # the pull up, https://github.com/apache/datafusion/issues/25519, and is not | ||
| # changed here: the three rows for `k = 3`, `k = 9` and `k = NULL` should be | ||
| # NULL. | ||
| query IIB | ||
| SELECT co.id, co.k, co.k IN (SELECT ci.k FROM ci WHERE ci.k = co.k GROUP BY ROLLUP(ci.k)) AS m FROM co ORDER BY k, id; | ||
| ---- | ||
| 1 1 true | ||
| 2 1 true | ||
| NULL 1 true | ||
| 2 2 true | ||
| NULL 3 false | ||
| 9 9 false | ||
| 5 NULL false |
There was a problem hiding this comment.
Result changed here, and the new value is the wrong one. Please read this one before the others.
The three rows for k = 3, k = 9 and k = NULL were NULL in the previous commit, and they are FALSE again now. FALSE is what main gives. DuckDB 1.5.2 gives NULL:
D SELECT co.id, co.k, co.k IN (SELECT ci.k FROM ci WHERE ci.k = co.k GROUP BY ROLLUP(ci.k)) AS m FROM co ORDER BY k, id;
┌───────┬───────┬───────┐
│ id │ k │ m │
│ ... │ ... │ ... │
│ NULL │ 3 │ NULL │
│ 9 │ 9 │ NULL │
│ 5 │ NULL │ NULL │
└───────┴───────┴───────┘
The NULL came from a guard that this commit removes on purpose. The guard cleared a flag for an outer join, a union or a grouping set, which made this one query correct without fixing its cause. The cause is in the pull up: it moves the correlated filter above the aggregate, and the grand-total row of ROLLUP is then one row for the whole table instead of one row for each outer row. The same plan gives a wrong EXISTS for the same tables, which the guard never touched:
# both main and this branch
SELECT co.id, co.k, EXISTS (SELECT 1 FROM ci WHERE ci.k = co.k GROUP BY ROLLUP(ci.k)) AS e FROM co;
-- gives false for k = 3, 9 and NULL; DuckDB gives true for every row
So the guard was not a fix for the bug, only for one of its symptoms, and keeping it means this rule carries a list of plan nodes that can put a NULL back into a column. I filed the cause as #25519 and left this query at the main result until that is fixed.
Tell me if you would rather keep the guard until #25519 lands. The change is one match arm, and the tests stay green either way.
| query I | ||
| SELECT id FROM naconst_corr_t1 WHERE 3 NOT IN (SELECT id FROM naconst_corr_t2 WHERE naconst_corr_t2.g > naconst_corr_t1.g); | ||
| ---- | ||
| 2 |
There was a problem hiding this comment.
This query gave a planning error before, and it now gives a result. DuckDB 1.5.2 agrees with the result.
The expectation was:
query error DataFusion error: Error during planning: null_aware LeftAnti join requires equi-join keys, but the join has none
3 is a constant, so the IN equality is not an equi-join key, and the correlation g > g is a residual filter. The rule asked for a null-aware anti join that the planner cannot build, so the query failed. It now becomes the mark joins that materialize the three-valued result, which need no equi-join key.
D SELECT a.id, a.g, (SELECT list(b.id) FROM naconst_corr_t2 b WHERE b.g > a.g) AS sub FROM naconst_corr_t1 a ORDER BY a.id;
┌───────┬───────┬──────────┐
│ id │ g │ sub │
├───────┼───────┼──────────┤
│ 1 │ 1 │ [NULL] │
│ 2 │ 2 │ NULL │
└───────┴───────┴──────────┘
D SELECT id FROM naconst_corr_t1 WHERE 3 NOT IN (SELECT id FROM naconst_corr_t2 WHERE naconst_corr_t2.g > naconst_corr_t1.g);
┌────┐
│ id │
├────┤
│ 2 │
└────┘
For id = 1 the subquery gives {NULL}, so 3 NOT IN ({NULL}) is UNKNOWN and the row goes. For id = 2 the subquery is empty, so the answer is TRUE and the row stays.
| 01)Projection: join_t1.t1_id, join_t1.t1_name, join_t1.t1_int | ||
| 02)--Filter: NOT CASE WHEN __correlated_sq_2.mark THEN Boolean(true) WHEN __correlated_sq_3.mark OR CAST(join_t1.t1_id AS Int64) + Int64(12) IS NULL AND __correlated_sq_4.mark THEN Boolean(NULL) ELSE Boolean(false) END | ||
| 03)----LeftMark Join: Filter: join_t1.t1_int > UInt32(0) | ||
| 04)------LeftMark Join: Filter: join_t1.t1_int > UInt32(0) | ||
| 05)--------LeftMark Join: CAST(join_t1.t1_id AS Int64) + Int64(12) = __correlated_sq_2.join_t2.t2_id + Int64(1) Filter: join_t1.t1_int > UInt32(0) |
There was a problem hiding this comment.
Only the plan changed here. The result of this query is the same before and after, and it is correct in both. I want to flag the cost, because this file tests the plan and not the result.
SELECT t1_id, t1_name, t1_int FROM join_t1 WHERE join_t1.t1_id + 12 NOT IN (SELECT join_t2.t2_id + 1 FROM join_t2 WHERE join_t1.t1_int > 0);
-- main, this branch and DuckDB 1.5.2 all give: 22 b 2
join_t2.t2_id holds no NULL in this file, so the old single anti join was correct for this data. The plan is now three joins for the same answer, and two of them are nested loop joins that carry only the outer-side predicate.
The reason is that the shape is not correct in general. The correlation join_t1.t1_int > 0 names no column of the subquery, so it stays as a residual filter, and the null-aware anti join executor ignores the residual when it decides whether a NULL makes the result UNKNOWN. Put one NULL in the subquery column and main drops every row:
CREATE TABLE jt1(t1_id INT, t1_int INT) AS VALUES (11,1),(22,2),(33,0);
CREATE TABLE jt2(t2_id INT) AS VALUES (23),(NULL);
SELECT t1_id, t1_int FROM jt1 WHERE jt1.t1_id + 12 NOT IN (SELECT jt2.t2_id + 1 FROM jt2 WHERE jt1.t1_int > 0);
-- main: (no rows)
-- this branch: 33 0
-- DuckDB: 33 0The row t1_id = 33 has t1_int = 0, so its correlated subquery is empty and NOT IN (<empty set>) is TRUE.
subquery_projection.slt now has this shape as a result test, with the tables oc and ic, so the gain is no longer only in this comment. #25336 is the executor gap behind it. When #25339 lands, this shape can go back to one anti join and this plan becomes small again. I can also hold this file at the old plan and let only the projected IN take the new path, if you would rather not pay three joins for a WHERE clause today.
| .all(|&e| can_pullup_over_aggregation(e)); | ||
| let (mut join_filters, subquery_filters) = | ||
| find_join_exprs(subquery_filter_exprs)?; | ||
| for expr in &join_filters { |
There was a problem hiding this comment.
Thank you, this is a real regression. I took your fix in 10ccc3e. correlated_filters now keeps the outer references, and filter_rejects_null matches each side of a conjunct against the side of the join that the key is on. Your two queries are in subquery_projection.slt (tables sh_o and sh_t), and DuckDB 1.5.2 gives the same results: NULL for the NULL row, and 8 for the NOT IN filter.
| // row", which is the weaker fact that this case needs, and which the first | ||
| // reading implies. `IS NULL` is two-valued on both sides, so neither test | ||
| // adds an UNKNOWN of its own. | ||
| let unknown_alias = alias.next("__correlated_sq"); |
There was a problem hiding this comment.
Agreed. I removed 50d6e4a from this PR. I also removed 1a9d842, which only fixed an error message that the merged join caused. The fallback for a residual filter is now the same plan as on main.
After the rebase onto #25339, a NOT IN filter with a residual filter no longer uses the fallback. It becomes one null-aware anti join, because that join now applies the residual filter. The nai_res_og plan that #25339 pins stays as it is, and the joins.slt plan is the same as on main again.
jayzhan211
left a comment
There was a problem hiding this comment.
Thanks @adriangb , here is a suggestion
| .all(|&e| can_pullup_over_aggregation(e)); | ||
| for expr in &subquery_filter_exprs { | ||
| if expr.contains_outer() && !self.correlated_filters.contains(expr) { | ||
| self.correlated_filters.push((*expr).clone()); |
There was a problem hiding this comment.
correlated_filters records a conjunct regardless of what sits above the Filter. A LEFT JOIN (filter on the nullable side) or a ROLLUP above it puts NULLs back into the key, but filter_rejects_null still reads u.y = o.x as "never NULL", so the NOT IN anti join loses null_aware. main returns the correct result for both queries below:
CREATE TABLE o(x INT) AS VALUES (1),(3),(NULL);
CREATE TABLE t(k INT) AS VALUES (1),(2);
CREATE TABLE u(y INT, k INT) AS VALUES (1,1);
SELECT x FROM o WHERE x NOT IN (
SELECT u.y FROM t LEFT JOIN (SELECT * FROM u WHERE u.y = o.x) AS u ON t.k = u.k);
-- expected (and main): no rows; this PR: 3, NULL
-- co/ci from subquery_projection.slt
SELECT co.id, co.k FROM co WHERE co.k NOT IN (
SELECT ci.k FROM ci WHERE ci.k = co.k GROUP BY ROLLUP(ci.k));
-- expected (and main): no rows; this PR: (NULL,3), (9,9), (5,NULL)Fix: drop the recorded filters when the pull up passes such a node, and add both queries to the slt:
fn f_up(&mut self, plan: LogicalPlan) -> Result<Transformed<LogicalPlan>> {
+ // A node that null-extends or regroups rows can put a NULL back into
+ // a column that a correlated filter below it rejected.
+ if may_reintroduce_nulls(&plan) {
+ self.correlated_filters.clear();
+ }
let subquery_schema = plan.schema();fn may_reintroduce_nulls(plan: &LogicalPlan) -> bool {
match plan {
LogicalPlan::Join(join) => matches!(
join.join_type,
JoinType::Left | JoinType::Right | JoinType::Full
),
LogicalPlan::Union(_) => true,
LogicalPlan::Aggregate(aggregate) => aggregate
.group_expr
.iter()
.any(|e| matches!(e, Expr::GroupingSet(_))),
_ => false,
}
}5148cd3 to
61a61ef
Compare
jayzhan211
left a comment
There was a problem hiding this comment.
Thanks @adriangb , LGTM!
kosiew
left a comment
There was a problem hiding this comment.
Thanks for the update. The projected IN / NOT IN rewrite now preserves three-valued NULL semantics across the covered expression, correlation, empty-subquery, grouping-set, and name-shadowing cases. The added plan and execution coverage addresses the earlier review points.
|
@kosiew @jayzhan211 thanks so much for helping get this across the line!! |
…ections A projected `IN` / `NOT IN` subquery used three mark joins to materialize its three-valued result. When the mark column of the first join is already exact, that join alone is the full rewrite: - no side of the `IN` predicate can be NULL in the scope of an outer row, or - the join is null-aware and the `IN` equality is its first key. The null-aware hash join also applies a residual non-equality filter when it decides whether a NULL makes the mark UNKNOWN (apache#25339). Key nullability is read from the key expressions, not only their columns, so `NULLIF(id, 1)` or `TRY_CAST(s AS INT)` over a non-nullable column is treated as nullable. A correlated conjunct that rejects NULL on a key, such as a correlation that repeats the `IN` predicate, keeps that key out of the null-aware decision (apache#25480). A `NOT IN` filter whose `IN` equality cannot be the first key of a null-aware anti join now falls back to the mark-join materialization instead of planning a wrong or unsupported anti join. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… a NULL
A correlated filter that rejects NULL on a key only bounds the subquery
result if nothing above it can put a NULL back. An outer join, a union or a
grouping set can. With `ROLLUP`, the grand-total row is a NULL for every
outer row, so a miss is UNKNOWN, but the join was planned without
null-aware semantics and gave FALSE:
SELECT k, k IN (SELECT i.k FROM i WHERE i.k = o.k GROUP BY ROLLUP(i.k))
FROM o;
The pull up now clears the recorded filters when it passes such a node.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The correlation on `<` now stays a residual filter of the one null-aware mark join, so the query no longer keeps the nested-loop plan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
8aa0b87 to
210d4e5
Compare
…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? Closes apache#25473. Closes apache#25474. ## Rationale for this change Issues apache#25473 and apache#25474 reported wrong results for `NOT IN (subquery)` when the subquery column is `NOT NULL` but the value being compared can evaluate to `NULL`. This includes NULL constants such as `CAST(NULL AS INT)` and nullable expressions over `NOT NULL` columns, such as a `CASE` expression without an `ELSE`. The fix itself landed in apache#25338, which determines null-awareness from the full `IN` operand expressions. This PR adds the missing regression coverage for both issues so this behavior does not regress. ## What changes are included in this PR? This is a test-only change. There are no production code or public API changes. It adds optimizer unit tests that verify: - A NULL constant compared against a non-nullable subquery column produces a projected `__correlated_sq_1_value` and a `null_aware` `LeftAnti` join. - A non-nullable expression such as `test.c + 1`, where `test.c` is non-nullable, continues to use a regular `LeftAnti` join without `null_aware`. It adds SQL logic tests for null-aware anti joins covering: - `CAST(NULL AS INT)`, untyped `NULL`, and `NULLIF(1, 1)`. - Nullable expressions over `NOT NULL` columns. - Nullable expressions on the subquery side. - NULL values against an empty subquery. - A non-NULL constant control. - A non-nullable `x + 1` control. - Logical and physical `EXPLAIN` plans for the NULL constant, nullable `CASE`, and non-nullable `x + 1` cases. - The NULL constant and nullable `CASE` cases with `datafusion.execution.target_partitions = 1`. It also adds SQL logic tests for the mark-join path, including the `... NOT IN (...) OR x = 99` form and a `SELECT`-list control that verifies the mark evaluates to `NULL`. ## Are these changes tested? Yes. This PR adds regression tests in: - `datafusion/optimizer/src/decorrelate_predicate_subquery.rs` - `datafusion/sqllogictest/test_files/null_aware_anti_join.slt` - `datafusion/sqllogictest/test_files/null_aware_mark_join.slt` The new tests verify both result correctness and expected logical/physical plan shapes, including that nullable operands use null-aware joins while non-nullable operands remain non-null-aware. Verification: ```bash cargo test -p datafusion-optimizer decorrelate_predicate_subquery cargo test -p datafusion-sqllogictest --test sqllogictests -- null_aware cargo test -p datafusion-sqllogictest --test sqllogictests -- subquery_projection ``` ## Are there any user-facing changes? No new user-facing behavior is introduced by this PR. The underlying correctness fix already landed in apache#25338. This PR only adds regression coverage for that behavior and does not change any public APIs. ## LLM-generated code disclosure This PR includes LLM-generated code and comments. All LLM-generated content has been manually reviewed.
Which issue does this PR close?
This PR supersedes #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 #24972, which added the decorrelation of
INsubqueries in a projection.Rationale for this change
An
INorNOT INsubquery 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.COALESCEduplicates 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 #25339.
INin the SELECT listCOALESCE((x IN (...))::boolean, false)INwith an equality predicateINsubqueries in separate columnsEXISTSINwith a non-equality predicateThe table is one run of the script below after the rebase onto
mainat39ca2c74a9,datafusion-clibuilt with--profile release-nonltoand the defaulttarget_partitions. Q6 in 3 more interleaved runs:main2.57 s to 3.18 s, this PR 0.022 s to 0.058 s.These shapes are the
projection_subquerybenchmark suite from #25346.Reproduction with datafusion-cli: script, plans and timings for each shape
Run the script with
datafusion-cli -f mre.sql. It creates two tables with 200000 rows each, then for each shape it printsEXPLAINand runs the query.datafusion-cliprints 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
mainat 22651d2 andtarget_partitions = 4. The plans of this PR for those shapes did not change with the rebase. Q6 was captured again after the rebase.Q1: Bare `IN` in the SELECT list, plans on main and on this PR
main, query time 190.730 s
this PR, query time 0.034 s
Q2: `COALESCE((x IN (...))::boolean, false)`, plans on main and on this PR
main, query time 351.579 s
this PR, query time 0.054 s
Q3: Correlated `IN` with an equality predicate, plans on main and on this PR
main, query time 1.372 s
this PR, query time 0.091 s
Q4: Two `IN` subqueries in separate columns, plans on main and on this PR
main, query time 296.676 s
this PR, query time 0.063 s
Q5: Correlated `EXISTS`, plans on main and on this PR
main, query time 0.021 s
this PR, query time 0.021 s
Q6: Correlated `IN` with a non-equality predicate, plans on main and on this PR
main (
39ca2c74a9), query time 2.489 sthis PR, query time 0.018 s
What changes are included in this PR?
build_joinnow reports whether the mark column of theLeftMarkjoin is exact under three-valued logic. It is exact when no side of theINpredicate can be NULL in scope, or when the join is null-aware and theINequality 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 (fix: correlated NOT IN with a non-equality correlation returns wrong results #25339).in_subquery_value_mark_joinbuilds this join first. If the mark is exact, it returns the mark column alone, orNOT markforNOT IN. The three join materialization stays for the other cases.The rule decides null-awareness from one question: can the
INvalue or the subquery output be NULL inside the scope of an outer row? Only then canINbe UNKNOWN.NULLIF(id, 1)orTRY_CAST(s AS INT)can be NULL over aNOT NULLcolumn.PullUpCorrelatedExprnow records every correlated conjunct, before it drops the conjunct that repeats theINpredicate. A conjunct such asy = x,y > xory IS NOT NULLis never TRUE for a NULLy, so it keeps a NULLyout 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.ROLLUP, the grand-total row is a NULL for every outer row.A
NOT INfilter uses the same scope-aware nullability for its null-awareLeftAntijoin. If theINequality does not become the first key of that join, the rule gives up and theNOT INbecomes the mark joins, which give the correct result.The regression test for fix: scan subqueries when advancing the extracted-alias generator #24574 in
projection_pushdown.sltnow uses a correlated subquery withLIMIT 1. The subquery then still reachesExtractLeafExpressions, 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.New sqllogictest cases in
subquery_projection.slt:EXPLAINguards: one hash mark join for each subquery, one null-aware mark join with a residual filter, a null-awareLeftAntijoin with two keys, and with one key and a residual filter.INandNOT INwith a NULL in the subquery, a NULL outer value,NOTinsideCASE, a projection over an aggregate, theCOALESCEand cast shape, and the residual filter shape.NOT NULLcolumns: aNULLIFvalue, aTRY_CASTvalue and aNULLIFsubquery output, in a projection and in aWHEREclause.INpredicate, and a subquery column with the same qualified name as the outer value.ROLLUP, as a projectedINwith anEXPLAINguard and as aNOT INfilter.The expected results agree with DuckDB 1.5.2, and the earlier cases also with PostgreSQL.
The
projection_subquerybenchmark docs no longer describe q07 as the nested-loop control.DuckDB comparison
For an absolute reference, the same seven queries in
datafusion-cliand in DuckDB 1.5.2, on the same two tables, release build, Apple M4 Pro.mainis the merge-base64871d9and "this PR" is3af87370c1. The later commits of this PR change the correlated andNOT INfilter paths, and, after the rebase onto #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.mainINCOALESCEoverININ, equalityINcolumnsNOT INEXISTSIN, non-equalityAll 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,
mainis 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:mainat39ca2c74a978 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:
maintakes 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.How to run it
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, becausegenerate_seriesgives a column of that name there and a column namedvaluein DataFusion:datafusion-clireports the time of each statement asElapsed, andduckdbreports it asRun Time (s): realafter.timer on.What is the testing strategy for this PR?
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 theINpredicate, and the same correlation below aROLLUP. The last one fails without the clearing in item 2.subquery_projection.sltlisted above.Are there any user-facing changes?
Plans for projected
INandNOT INsubqueries 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:
INorNOT INwhose key is an expression that can be NULL over aNOT NULLcolumn, such asNULLIF(id, 1). For a projectedIN,mainreturnsfalseinstead ofNULL. For aNOT INfilter,mainkeeps rows that it must drop. The correlated form of that is Wrong results: correlatedNOT IN (subquery)ignores NULLs in the subquery #25347.NOT INwhose correlation repeats theINpredicate,x NOT IN (SELECT y FROM t WHERE y = x).maindrops rows that it must keep (Wrong results: correlatedNOT INloses the correlation when it is an equality on theINvalue column #25480).INwhose correlation repeats theINpredicate below aROLLUP,k IN (SELECT i.k FROM i WHERE i.k = o.k GROUP BY ROLLUP(i.k)).mainfails to plan it since fix: keep a correlated filter below an aggregate with a grouping set #25529. It now gives NULL for a miss, as DuckDB does.🤖 Generated with Claude Code