apache / apache/datafusion

Decorrelate scalar subqueries with more complex filter expressions

Open
#14,554 15 comments 0 reactions 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Is your feature request related to a problem or challenge?

Datafusion already support decorrelating simple scalar subqueries in this PR: https://github.com/apache/datafusion/pull/6457

This follow the first approach in [TUM paper](https://btw-2015.informatik.uni-hamburg.de/res/proceedings/Hauptband/Wiss/Neumann-Unnesting_Arbitrary_Querie.pdf) (simple unnesting), and allow decorrelating [this simple query](https://github.com/apache/datafusion/blob/873d4178ab3c98ca6f78ef553c8e33e12dd342fb/datafusion/core/tests/sqllogictests/test_files/subquery.slt#L797)

```
explain select t1.t1_int from t1 where (select count(*) from t2 where t1.t1_id = t2.t2_id) < t1.t1_int
```
However, if we add an `or` condition this subquery
```
explain select t1.t1_int from t1 where (select count(*) from t2 where t1.t1_id = t2.t2_id or t1.t1_name=t2.t2_name) < t1.t1_int
```

Datafusion cannot decorrelate it
```
+--------------+----------------------------------------------------------------------------------------+
| plan_type | plan |
+--------------+----------------------------------------------------------------------------------------+
| logical_plan | Projection: t1.t1_int |
| | Filter: () < CAST(t1.t1_int AS Int64) |
| | Subquery: |
| | Projection: count(*) |
| | Aggregate: groupBy=[[]], aggr=[[count(Int64(1)) AS count(*)]] |
| | Filter: outer_ref(t1.t1_id) = t2.t2_id OR outer_ref(t1.t1_name) = t2.t2_name |
| | TableScan: t2 |
| | TableScan: t1 projection=[t1_id, t1_name, t1_int] |
+--------------+----------------------------------------------------------------------------------------+
```

### Describe the solution you'd like

Support decorrelating this query following the second method mentioned in the paper

### Describe alternatives you've considered

_No response_

### Additional context

General framework for decorrelation maybe discussed here https://github.com/apache/datafusion/issues/5492

But the steps needed to make this work is followed

Allow decorrelation for this type of filter exprs in this code: https://github.com/apache/datafusion/blob/813220d54f08c5203ad79bfb066ca638abe208ed/datafusion/optimizer/src/decorrelate.rs#L162

Add more logic to handle complex query decorrelation:
- Build domain/magic relation
- Rewrite the subquery to join inner table (table of the subquery) with domain/magic relation using its complex filter expression (i.e `t2.t2_id = domain.t1_id OR t2.t2_name = domain.t1_name`)
- Rewrite aggregation to group by the additional columns mentioned in the domain/magic relation
- Join the outer relation with the newly built aggregation

For example the above mentioned query may be rewritten like
```
explain select t1.t1_int from t1,
(
select count(*) as count_all, domain.t1_id as t1_id, domain.t1_name as t1_name from (
select distinct t1_id, t1_name from t1
) as domain join t2 where t2.t2_id = domain.t1_id or t2.t2_name=domain.t1_name
group by domain.t1_id, domain.t1_name
) as pulled_up
where t1.t1_id=pulled_up.t1_id and pulled_up.count_all < t1.t1_int
```

Logical plan may look like
```
| logical_plan | Projection: t1.t1_int |
| | Inner Join: t1.t1_id = pulled_up.t1_id Filter: pulled_up.count_all < CAST(t1.t1_int AS Int64) |
| | TableScan: t1 projection=[t1_id, t1_int] |
| | SubqueryAlias: pulled_up |
| | Projection: count(*) AS count_all, domain.t1_id |
| | Aggregate: groupBy=[[domain.t1_id, domain.t1_name]], aggr=[[count(Int64(1)) AS count(*)]] |
| | Projection: domain.t1_id, domain.t1_name |
| | Inner Join: Filter: t2.t2_id = domain.t1_id OR t2.t2_name = domain.t1_name |
| | SubqueryAlias: domain |
| | Aggregate: groupBy=[[t1.t1_id, t1.t1_name]], aggr=[[]] |
| | TableScan: t1 projection=[t1_id, t1_name] |
| | TableScan: t2 projection=[t2_id, t2_name]
```

Contributor guide

Open the contributing guide

Research direction

Start in datafusion/optimizer/src/decorrelate.rs around the filter-expression handling at line 162, then review the existing scalar-subquery case and the subquery example in datafusion/core/tests/sqllogictests/test_files/subquery.slt around line 797. Implement the domain/magic-relation, join, grouping, and aggregation steps described in the issue, then run the relevant sqllogictest and confirm the complex OR-filter query is decorrelated.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust, sql
Domain
databases
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.