Repository navigation
Introduce dependent join LogicalPlan to support complex subquery decorrelation - #21555
Closed
duongcongtoai wants to merge 10 commits into
Closed
duongcongtoai wants to merge 10 commits into
duongcongtoai wants to merge 10 commits into
Conversation
LogicalPlan to support complex subquery decorrelation
This was referenced Apr 23, 2026
…-dependent-join-integration
Contributor
|
@neilconway or @mbutrovich I wonder if you are able to help review this one (it is pretty complex, but heads us down the path to proper subquery support) |
duongcongtoai
marked this pull request as ready for review
May 4, 2026 13:10
|
Thank you for your contribution. Unfortunately, this pull request is stale because it has been open 60 days with no activity. Please remove the stale label or comment or this will be closed in 7 days. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Which issue does this PR close?
#16059 has completed, but the result is not persisted into datafusion.
=> This PR brings back all the changes inside the GSOC work and a complete POC.
I'll try to break it down into smaller components and bring them into datafusion:
Prerequisite
Major part
Rationale for this change
1. 📚 Academic Background & The Paper's Core Concept
The Problem
Traditional subquery decorrelation rules in query optimizers typically rely on simple pattern-matching heuristics. They successfully unnest basic
IN,EXISTS, and simple scalar subqueries with equality predicates. However, they struggle or fail entirely when encountering complex correlations, such as:OR).HAVINGclause conditions).The Solution: Thomas Neumann's Arbitrary Unnesting Paper
This PR implements the groundwork for a systematic, algebraic approach to unnesting any arbitrary correlated subquery, based on:
Instead of converting subqueries directly to joins using heuristic rules, the paper abstracts the correlation boundary by introducing the Dependent Join ($\bowtie_d$ ) operator. A Dependent Join is a join where the evaluation of the right-hand side (RHS) depends on values produced by the left-hand side (LHS).
The 3-Phase Architecture
As per the paper:
The complete unnesting framework is mapped into DataFusion as:
Detects all outer-column references (accessors) and their corresponding tables/operators (providers). It computes their Lowest Common Ancestor (LCA), extracts the subquery, and inserts a
DependentJoinoperator at the LCA.Applies a top-down decorrelation algorithm to eliminate the
DependentJoinnodes, turning them into standard physical joins (e.g., Left, Semi, Anti Joins).Simplifies the resulting query tree, merging redundant operators and removing unnecessary join/scan boundaries via Delimination scans (DelimGets).
2. 🔍 How Phase 1 Works: The Stack-Based LCA Algorithm
To place a
DependentJoinoperator correctly, the planner must find the exact boundary where the correlation is introduced.DependentJoinmust be placed here.Stack-Based LCA Tracking
The paper suggests using an indexed algebra to do LCA computations in$O(\log n)$ . If unsupported, it notes:
To implement this annotation similar to the paper, in DataFusion we use the tree traversal API on the root
LogicalPlannode, specifically methodrewrite_with_subqueries. We keep track of the path (thestackof node IDs) from the root to the current node.Crucial Traversal Property: When evaluating an operator, its subquery expressions are always visited first before its standard relational inputs. This guarantees we encounter and record the accessor ($A$ ) before we reach the provider ($P$ ).
Step-by-Step Traversal Trace
Given the paper's example query:
Let's trace this logic using the exact node numbers from the diagram below:
[1]. It contains a subquery, so we visit its subquery branch first:[2] Subquery.outer_ref(customer.c_custkey). We recordc_custkeyintounresolved_outer_ref_columnswith the accessor stack:[1, 2, 3, 4].[5] Subquery.orders.o_orderkeyandcustomer.c_custkeyintounresolved_outer_ref_columnswith the accessor stack:[1, 2, 3, 4, 5, 6, 7].[8] TableScan: lineitemand unwind).[4], we visit its relational input[9] TableScan: orders.[9]is the provider fororders.o_orderkey. Its stack is[1, 2, 3, 4, 9].[1, 2, 3, 4, 5, 6, 7]foro_orderkey.[1, 2, 3, 4, 9]and the accessor's stack[1, 2, 3, 4, 5, 6, 7]. The deepest common prefix is Node[4]. We annotate Node[4].o_orderkeyat this level have been resolved. Node[4]is rewritten into aDependentJoinwhere the RHS is the subquery.[1], we visit its relational input branch reaching[11] TableScan: customer.[11]is the provider forcustomer.c_custkey. Its stack is[1, 10, 11].c_custkey:[1, 2, 3, 4]and[1, 2, 3, 4, 5, 6, 7].[1]. We annotate Node[1].c_custkeyhave been resolved. Node[1]is rewritten into aDependentJoinwhere the RHS is the top-level subquery.3. 🌳 Logical Plan Tree Restructuring
Before Rewrite:
After Rewrite (Dependent Joins inserted at LCAs):
4. 📐 Key Data Structures
Here are the primary structs introduced in this PR to represent dependency annotations and power the optimization passes:
LogicalPlan::DependentJoinThis is the new logical operator that formally models the dependent join.
schemaDFSchemaRefcorrelated_columnsVec<CorrelatedColumnInfo>subquery_exprOption<Expr>subquery_depthusizeleftArc<LogicalPlan>rightArc<LogicalPlan>subquery_nameStringlateral_join_conditionOption<(JoinType, Expr)>any_joinboolANY/SOMEorEXISTS/INsubquery comparison expression.CorrelatedColumnInfoStores specific metadata for columns that cross the correlation boundary.
DependentJoinRewriterThe stateful tree rewriter that implements
TreeNodeRewriterto execute Phase 1.ColumnAccessTracks an active column access before it is resolved by its provider.
5. 🚀 What’s Next (Roadmap)
The collections of nodes that provides the columns will be persisted and passed to the next round of decorrelation optimization (to construct delim_get).
DependentJoinDecorrelator): We will introduce a top-down decorrelation algorithm to eliminateDependentJoinnodes, turning them into standard physical joins (e.g., Left, Semi, Anti Joins).Deliminator): We will translate the remaining annotations and domain providers into physical/logical Delimination Scans (DelimGets) to yield a fully decorrelated, high-performance query plan.