Skip to content

Commit 6a0ce22

Browse files
committed
fix(query): match nullable scalar correlation keys
1 parent 9661185 commit 6a0ce22

6 files changed

Lines changed: 283 additions & 214 deletions

File tree

src/query/sql/src/planner/optimizer/optimizers/operator/decorrelate/decorrelate.rs

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -275,6 +275,20 @@ impl SubqueryDecorrelatorOptimizer {
275275
&mut left_conditions,
276276
)?;
277277

278+
// These conditions reconnect the flattened subquery result to the outer row.
279+
// They are internal correlation keys rather than user-written equality
280+
// predicates, so NULL correlation groups must match each other.
281+
let mut is_null_equal = Vec::new();
282+
for (i, (l, r)) in left_conditions
283+
.iter()
284+
.zip(right_conditions.iter())
285+
.enumerate()
286+
{
287+
if l.data_type().is_nullable() || r.data_type().is_nullable() {
288+
is_null_equal.push(i);
289+
}
290+
}
291+
278292
let join_type = if matches!(subquery.contain_agg, Some(true)) && {
279293
let rel_expr = RelExpr::with_s_expr(&subquery.subquery);
280294
rel_expr
@@ -292,7 +306,7 @@ impl SubqueryDecorrelatorOptimizer {
292306
equi_conditions: JoinEquiCondition::new_conditions(
293307
left_conditions,
294308
right_conditions,
295-
vec![],
309+
is_null_equal,
296310
),
297311
non_equi_conditions: vec![],
298312
join_type,

src/query/sql/test-support/data/results/tpcds/Q01_optimized.txt

Lines changed: 70 additions & 70 deletions
Original file line numberDiff line numberDiff line change
@@ -15,80 +15,80 @@ TopN
1515
├── other filters: []
1616
├── Exchange(Broadcast)
1717
│ └── Join(Inner)
18-
│ ├── build keys: [store_returns.sr_store_sk (#103)]
18+
│ ├── build keys: [store.s_store_sk (#49)]
1919
│ ├── probe keys: [store_returns.sr_store_sk (#7)]
20-
│ ├── other filters: [gt(Sum(sr_return_amt) (#48), sum(ctr_total_return) / if(count(ctr_total_return) = 0, 1, count(ctr_total_return)) * 1.2 (#147))]
20+
│ ├── other filters: []
2121
│ ├── Exchange(Broadcast)
22-
│ │ └── Join(Inner)
23-
│ │ ├── build keys: [store_returns.sr_store_sk (#103)]
24-
│ │ ├── probe keys: [store.s_store_sk (#49)]
25-
│ │ ├── other filters: []
26-
│ │ ├── Exchange(Broadcast)
27-
│ │ │ └── EvalScalar
28-
│ │ │ ├── scalars: [store_returns.sr_store_sk (#103) AS (#103), multiply(divide(sum(ctr_total_return) (#145), if(eq(count(ctr_total_return) (#146), 0), 1, count(ctr_total_return) (#146))), 1.2) AS (#147)]
29-
│ │ │ └── Aggregate(Final)
30-
│ │ │ ├── group items: [store_returns.sr_store_sk (#103) AS (#103)]
31-
│ │ │ ├── aggregate functions: [sum(Sum(sr_return_amt) (#144)) AS (#145), count(Sum(sr_return_amt) (#144)) AS (#146)]
32-
│ │ │ └── Aggregate(Partial)
33-
│ │ │ ├── group items: [store_returns.sr_store_sk (#103) AS (#103)]
34-
│ │ │ ├── aggregate functions: [sum(Sum(sr_return_amt) (#144)) AS (#145), count(Sum(sr_return_amt) (#144)) AS (#146)]
35-
│ │ │ └── Exchange(Hash)
36-
│ │ │ ├── Exchange(Hash): keys: [store_returns.sr_store_sk (#103)]
37-
│ │ │ └── Aggregate(Final)
38-
│ │ │ ├── group items: [store_returns.sr_customer_sk (#99) AS (#99), store_returns.sr_store_sk (#103) AS (#103)]
39-
│ │ │ ├── aggregate functions: [sum(store_returns.sr_return_amt (#107)) AS (#144)]
40-
│ │ │ └── Aggregate(Partial)
41-
│ │ │ ├── group items: [store_returns.sr_customer_sk (#99) AS (#99), store_returns.sr_store_sk (#103) AS (#103)]
42-
│ │ │ ├── aggregate functions: [sum(store_returns.sr_return_amt (#107)) AS (#144)]
43-
│ │ │ └── Exchange(Hash)
44-
│ │ │ ├── Exchange(Hash): keys: [store_returns.sr_customer_sk (#99)]
45-
│ │ │ └── EvalScalar
46-
│ │ │ ├── scalars: [store_returns.sr_customer_sk (#99) AS (#99), store_returns.sr_store_sk (#103) AS (#103), store_returns.sr_return_amt (#107) AS (#107), store_returns.sr_returned_date_sk (#96) AS (#151), date_dim.d_date_sk (#116) AS (#152), date_dim.d_year (#122) AS (#153)]
47-
│ │ │ └── Join(Inner)
48-
│ │ │ ├── build keys: [date_dim.d_date_sk (#116)]
49-
│ │ │ ├── probe keys: [store_returns.sr_returned_date_sk (#96)]
50-
│ │ │ ├── other filters: []
51-
│ │ │ ├── Exchange(Broadcast)
52-
│ │ │ │ └── Scan
53-
│ │ │ │ ├── table: default.date_dim (#5)
54-
│ │ │ │ ├── filters: [eq(date_dim.d_year (#122), 2001)]
55-
│ │ │ │ ├── order by: []
56-
│ │ │ │ └── limit: NONE
57-
│ │ │ └── Scan
58-
│ │ │ ├── table: default.store_returns (#4)
59-
│ │ │ ├── filters: []
60-
│ │ │ ├── order by: []
61-
│ │ │ └── limit: NONE
62-
│ │ └── Scan
63-
│ │ ├── table: default.store (#2)
64-
│ │ ├── filters: [eq(store.s_state (#73), 'TN')]
65-
│ │ ├── order by: []
66-
│ │ └── limit: NONE
67-
│ └── Aggregate(Final)
68-
│ ├── group items: [store_returns.sr_customer_sk (#3) AS (#3), store_returns.sr_store_sk (#7) AS (#7)]
69-
│ ├── aggregate functions: [sum(store_returns.sr_return_amt (#11)) AS (#48)]
70-
│ └── Aggregate(Partial)
22+
│ │ └── Scan
23+
│ │ ├── table: default.store (#2)
24+
│ │ ├── filters: [eq(store.s_state (#73), 'TN')]
25+
│ │ ├── order by: []
26+
│ │ └── limit: NONE
27+
│ └── Join(Inner)
28+
│ ├── build keys: [store_returns.sr_store_sk (#103)]
29+
│ ├── probe keys: [store_returns.sr_store_sk (#7)]
30+
│ ├── other filters: [gt(Sum(sr_return_amt) (#48), sum(ctr_total_return) / if(count(ctr_total_return) = 0, 1, count(ctr_total_return)) * 1.2 (#147))]
31+
│ ├── Exchange(Broadcast)
32+
│ │ └── EvalScalar
33+
│ │ ├── scalars: [store_returns.sr_store_sk (#103) AS (#103), multiply(divide(sum(ctr_total_return) (#145), if(eq(count(ctr_total_return) (#146), 0), 1, count(ctr_total_return) (#146))), 1.2) AS (#147)]
34+
│ │ └── Aggregate(Final)
35+
│ │ ├── group items: [store_returns.sr_store_sk (#103) AS (#103)]
36+
│ │ ├── aggregate functions: [sum(Sum(sr_return_amt) (#144)) AS (#145), count(Sum(sr_return_amt) (#144)) AS (#146)]
37+
│ │ └── Aggregate(Partial)
38+
│ │ ├── group items: [store_returns.sr_store_sk (#103) AS (#103)]
39+
│ │ ├── aggregate functions: [sum(Sum(sr_return_amt) (#144)) AS (#145), count(Sum(sr_return_amt) (#144)) AS (#146)]
40+
│ │ └── Exchange(Hash)
41+
│ │ ├── Exchange(Hash): keys: [store_returns.sr_store_sk (#103)]
42+
│ │ └── Aggregate(Final)
43+
│ │ ├── group items: [store_returns.sr_customer_sk (#99) AS (#99), store_returns.sr_store_sk (#103) AS (#103)]
44+
│ │ ├── aggregate functions: [sum(store_returns.sr_return_amt (#107)) AS (#144)]
45+
│ │ └── Aggregate(Partial)
46+
│ │ ├── group items: [store_returns.sr_customer_sk (#99) AS (#99), store_returns.sr_store_sk (#103) AS (#103)]
47+
│ │ ├── aggregate functions: [sum(store_returns.sr_return_amt (#107)) AS (#144)]
48+
│ │ └── Exchange(Hash)
49+
│ │ ├── Exchange(Hash): keys: [store_returns.sr_customer_sk (#99)]
50+
│ │ └── EvalScalar
51+
│ │ ├── scalars: [store_returns.sr_customer_sk (#99) AS (#99), store_returns.sr_store_sk (#103) AS (#103), store_returns.sr_return_amt (#107) AS (#107), store_returns.sr_returned_date_sk (#96) AS (#151), date_dim.d_date_sk (#116) AS (#152), date_dim.d_year (#122) AS (#153)]
52+
│ │ └── Join(Inner)
53+
│ │ ├── build keys: [date_dim.d_date_sk (#116)]
54+
│ │ ├── probe keys: [store_returns.sr_returned_date_sk (#96)]
55+
│ │ ├── other filters: []
56+
│ │ ├── Exchange(Broadcast)
57+
│ │ │ └── Scan
58+
│ │ │ ├── table: default.date_dim (#5)
59+
│ │ │ ├── filters: [eq(date_dim.d_year (#122), 2001)]
60+
│ │ │ ├── order by: []
61+
│ │ │ └── limit: NONE
62+
│ │ └── Scan
63+
│ │ ├── table: default.store_returns (#4)
64+
│ │ ├── filters: []
65+
│ │ ├── order by: []
66+
│ │ └── limit: NONE
67+
│ └── Aggregate(Final)
7168
│ ├── group items: [store_returns.sr_customer_sk (#3) AS (#3), store_returns.sr_store_sk (#7) AS (#7)]
7269
│ ├── aggregate functions: [sum(store_returns.sr_return_amt (#11)) AS (#48)]
73-
│ └── Exchange(Hash)
74-
│ ├── Exchange(Hash): keys: [store_returns.sr_customer_sk (#3)]
75-
│ └── EvalScalar
76-
│ ├── scalars: [store_returns.sr_customer_sk (#3) AS (#3), store_returns.sr_store_sk (#7) AS (#7), store_returns.sr_return_amt (#11) AS (#11), store_returns.sr_returned_date_sk (#0) AS (#148), date_dim.d_date_sk (#20) AS (#149), date_dim.d_year (#26) AS (#150)]
77-
│ └── Join(Inner)
78-
│ ├── build keys: [date_dim.d_date_sk (#20)]
79-
│ ├── probe keys: [store_returns.sr_returned_date_sk (#0)]
80-
│ ├── other filters: []
81-
│ ├── Exchange(Broadcast)
82-
│ │ └── Scan
83-
│ │ ├── table: default.date_dim (#1)
84-
│ │ ├── filters: [eq(date_dim.d_year (#26), 2001)]
85-
│ │ ├── order by: []
86-
│ │ └── limit: NONE
87-
│ └── Scan
88-
│ ├── table: default.store_returns (#0)
89-
│ ├── filters: []
90-
│ ├── order by: []
91-
│ └── limit: NONE
70+
│ └── Aggregate(Partial)
71+
│ ├── group items: [store_returns.sr_customer_sk (#3) AS (#3), store_returns.sr_store_sk (#7) AS (#7)]
72+
│ ├── aggregate functions: [sum(store_returns.sr_return_amt (#11)) AS (#48)]
73+
│ └── Exchange(Hash)
74+
│ ├── Exchange(Hash): keys: [store_returns.sr_customer_sk (#3)]
75+
│ └── EvalScalar
76+
│ ├── scalars: [store_returns.sr_customer_sk (#3) AS (#3), store_returns.sr_store_sk (#7) AS (#7), store_returns.sr_return_amt (#11) AS (#11), store_returns.sr_returned_date_sk (#0) AS (#148), date_dim.d_date_sk (#20) AS (#149), date_dim.d_year (#26) AS (#150)]
77+
│ └── Join(Inner)
78+
│ ├── build keys: [date_dim.d_date_sk (#20)]
79+
│ ├── probe keys: [store_returns.sr_returned_date_sk (#0)]
80+
│ ├── other filters: []
81+
│ ├── Exchange(Broadcast)
82+
│ │ └── Scan
83+
│ │ ├── table: default.date_dim (#1)
84+
│ │ ├── filters: [eq(date_dim.d_year (#26), 2001)]
85+
│ │ ├── order by: []
86+
│ │ └── limit: NONE
87+
│ └── Scan
88+
│ ├── table: default.store_returns (#0)
89+
│ ├── filters: []
90+
│ ├── order by: []
91+
│ └── limit: NONE
9292
└── Scan
9393
├── table: default.customer (#3)
9494
├── filters: []

0 commit comments

Comments
 (0)