Skip to content

Introduce dependent join LogicalPlan to support complex subquery decorrelation - #21555

Closed
duongcongtoai wants to merge 10 commits into
apache:mainfrom
duongcongtoai:feat-dependent-join-integration
Closed

duongcongtoai wants to merge 10 commits into
apache:mainfrom
duongcongtoai:feat-dependent-join-integration

Conversation

@duongcongtoai

@duongcongtoai duongcongtoai commented Apr 11, 2026 •

Copy link
Copy Markdown
Contributor

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

  • logical/physical operator for delimget
  • left singlejoin support

Major part

  • implement DependentJoinRewriter <- This PR
  • implement DependentJoinDecorrelator
  • implement Deliminator

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:

  • Correlations nested deep inside disjunctions (OR).
  • Correlations inside projection lists or aggregates (e.g., HAVING clause conditions).
  • Multi-level deeply nested subqueries.

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:

"Unnesting Arbitrary Subqueries" (Thomas Neumann & Viktor Leis, BTW 2015 / Technical Report [Ne24]).

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:

We split the algorithm into three parts: First, a preparatory phase that identifies all non-trivial dependent joins and annotates them with information that the main algorithm needs. Second, the logic to eliminate dependent joins, which will be called for all non-trivial dependent joins in top-to-bottom order and which is the main algorithm, and third, the unnesting rules for individual operators.

The complete unnesting framework is mapped into DataFusion as:

  1. Phase 1: Preparatory Phase (LCA Identification & Annotation) — [THIS PR]
    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 DependentJoin operator at the LCA.
  2. Phase 2: Dependent Join Elimination (Decorrelation) — [Future PR]
    Applies a top-down decorrelation algorithm to eliminate the DependentJoin nodes, turning them into standard physical joins (e.g., Left, Semi, Anti Joins).
  3. Phase 3: Deliminator (Simplification) — [Future PR]
    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 DependentJoin operator correctly, the planner must find the exact boundary where the correlation is introduced.

  • Let $A$ be the operator/expression inside the subquery that accesses an outer column.
  • Let $P$ be the operator/table on the LHS that provides that column.
  • The Lowest Common Ancestor (LCA) of $A$ and $P$ represents the point in the tree where the correlation boundary is crossed. The DependentJoin must 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:

"the same information can be computed with worse asymptotic complexity by keeping track of the column sets that are available in the different parts of the tree."

To implement this annotation similar to the paper, in DataFusion we use the tree traversal API on the root LogicalPlan node, specifically method rewrite_with_subqueries. We keep track of the path (the stack of 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:

SELECT *
FROM customer
WHERE c_mktsegment='AUTOMOBILE' AND
    (SELECT COUNT(*) FROM orders
        WHERE o_custkey=c_custkey AND
            (SELECT SUM(l_extendedprice) FROM lineitem
                WHERE l_orderkey=o_orderkey
            )>300000
    )>5

Let's trace this logic using the exact node numbers from the diagram below:

  1. Down-Pass -> Filter [1]: We start descending at Node [1]. It contains a subquery, so we visit its subquery branch first: [2] Subquery.
  2. Down-Pass -> Filter [4]:
    • It contains an outer reference outer_ref(customer.c_custkey). We record c_custkey into unresolved_outer_ref_columns with the accessor stack: [1, 2, 3, 4].
    • It also contains a nested subquery expression, so we continue the traversal toward this subquery: [5] Subquery.
  3. Down-Pass -> Filter [7]:
    • It contains two outer references. We record orders.o_orderkey and customer.c_custkey into unresolved_outer_ref_columns with the accessor stack: [1, 2, 3, 4, 5, 6, 7].
    • (We hit leaf node [8] TableScan: lineitem and unwind).
  4. Down-Pass -> TableScan [9]: Unwinding back to [4], we visit its relational input [9] TableScan: orders.
    • Node [9] is the provider for orders.o_orderkey. Its stack is [1, 2, 3, 4, 9].
    • We check our unresolved accessors and find [1, 2, 3, 4, 5, 6, 7] for o_orderkey.
    • LCA Computation: We compute the LCA of the provider's stack [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].
  5. Up-Pass -> Filter [4]: All outer references mapping to o_orderkey at this level have been resolved. Node [4] is rewritten into a DependentJoin where the RHS is the subquery.
  6. Down-Pass -> TableScan [11]: Unwinding back to the top-level [1], we visit its relational input branch reaching [11] TableScan: customer.
    • Node [11] is the provider for customer.c_custkey. Its stack is [1, 10, 11].
    • We have two unresolved accessors for c_custkey: [1, 2, 3, 4] and [1, 2, 3, 4, 5, 6, 7].
    • LCA Computation: We compute the LCA of the provider's stack and the accessor stacks. The deepest common prefix for all paths is Node [1]. We annotate Node [1].
  7. Up-Pass -> Filter [1]: All outer references for c_custkey have been resolved. Node [1] is rewritten into a DependentJoin where the RHS is the top-level subquery.

3. 🌳 Logical Plan Tree Restructuring

Before Rewrite:

                                      ┌─────────────────────────┐
                                      │ [1] Filter              │
                                      │ (<subquery>) > Int32(5) │
                                      └───────────┬─────────────┘
                   ┌──────────────────────────────┴──────────────────────────────┐
                   ▼                                                             ▼
             ┌──────────────┐                                ┌────────────────────────────────────────────┐
             │ [2] Subquery │                                │ [10] Filter                                │
             └──────┬───────┘                                │ customer.c_mktsegment = Utf8("AUTOMOBILE") │
                    │                                        └───────────────────────┬────────────────────┘
                    ▼                                                                ▼
             ┌───────────────┐                                               ┌────────────────┐
             │ [3] Aggregate │                                               │ [11] TableScan │
             │ count(*)      │                                               │ customer       │
             └──────┬────────┘                                               └────────────────┘
                    │
                    ▼
┌──────────────────────────────────────────────────────────────────────────────────┐
│ [4] Filter                                                                       │
│ orders.o_custkey = outer_ref(customer.c_custkey) AND (<subquery>) > Int32(30000) │
└───────────────────┬──────────────────────────────────────────────────────────────┘
         ┌──────────┴──────────┐
         ▼                     ▼
  ┌──────────────┐      ┌───────────────┐
  │ [5] Subquery │      │ [9] TableScan │
  └──────┬───────┘      │ orders        │
         │              └───────────────┘
         ▼
┌─────────────────────────────────┐
│ [6] Aggregate                   │
│ count(lineitem.l_extendedprice) │
└────────────────┬────────────────┘
                 ▼
┌───────────────────────────────────────────────────────────────────────────────────────────────────────────┐
│ [7] Filter                                                                                                │
│ lineitem.l_orderkey = outer_ref(orders.o_orderkey) AND lineitem.l_custkey = outer_ref(customer.c_custkey) │
└────────────────────────┬──────────────────────────────────────────────────────────────────────────────────┘
                         ▼
                 ┌───────────────┐
                 │ [8] TableScan │
                 │ lineitem      │
                 └───────────────┘

After Rewrite (Dependent Joins inserted at LCAs):

                                       ┌────────────────┐
                                       │ [1] Projection │
                                       └───────┬────────┘
                                               │
                                               ▼
                                     ┌───────────────────┐
                                     │ [2] Filter        │
                                     │ __scalar_sq_2 > 5 │
                                     └─────────┬─────────┘
                                               │
                                               ▼
                                 ┌───────────────────────┐
                                 │ [3] DependentJoin     │
                                 │ on customer.c_custkey │
                                 └─────────────┬─────────┘
                       ┌───────────────────────┴───────────────────────────────────────┐
                       ▼                                                               ▼
               ┌───────────────┐                                       ┌────────────────────────────────────────────┐
               │ [4] Aggregate │                                       │ [13] Filter                                │
               │ count(*)      │                                       │ customer.c_mktsegment = Utf8("AUTOMOBILE") │
               └───────┬───────┘                                       └───────────────────────┬────────────────────┘
                       │                                                                       │
                       ▼                                                                       ▼
               ┌────────────────┐                                                      ┌────────────────┐
               │ [5] Projection │                                                      │ [14] TableScan │
               └───────┬────────┘                                                      │ customer       │
                       │                                                               └────────────────┘
                       ▼
             ┌───────────────────────┐
             │ [6] Filter            │
             │ __scalar_sq_1 > 30000 │
             └─────────┬─────────────┘
                       │
                       ▼
             ┌──────────────────────┐
             │ [7] DependentJoin    │
             │ on orders.o_orderkey │
             └─────────┬────────────┘
      ┌────────────────┴─────────────────────────────┐
      ▼                                              ▼
┌─────────────────────────────────┐        ┌──────────────────────────────────────────────────┐
│ [8] Aggregate                   │        │ [11] Filter                                      │
│ count(lineitem.l_extendedprice) │        │ orders.o_custkey = outer_ref(customer.c_custkey) │
└─────────────┬───────────────────┘        └────────────────────────┬─────────────────────────┘
              │                                                     │
              ▼                                                     ▼
┌────────────────────────────────────────────────────────┐  ┌────────────────┐
│ [9] Filter                                             │  │ [12] TableScan │
│ lineitem.l_orderkey = outer_ref(orders.o_orderkey) AND │  │ orders         │
│ lineitem.l_custkey = outer_ref(customer.c_custkey)     │  └────────────────┘
└─────────────────────────────┬──────────────────────────┘
                              │
                              ▼
                      ┌────────────────┐
                      │ [10] TableScan │
                      │ lineitem       │
                      └────────────────┘

4. 📐 Key Data Structures

Here are the primary structs introduced in this PR to represent dependency annotations and power the optimization passes:

LogicalPlan::DependentJoin

This is the new logical operator that formally models the dependent join.

Field Type Description
schema DFSchemaRef The output schema of the join.
correlated_columns Vec<CorrelatedColumnInfo> Outer references consumed by the RHS that are produced by the LHS.
subquery_expr Option<Expr> The original subquery expression.
subquery_depth usize Correlation depth (begins at 1).
left Arc<LogicalPlan> The LHS input, supplying correlated columns.
right Arc<LogicalPlan> The RHS input, containing the subquery tree accessing the LHS.
subquery_name String Stable name/alias assigned to this subquery.
lateral_join_condition Option<(JoinType, Expr)> Custom condition when representing a lateral join.
any_join bool True if there exists an ANY/SOME or EXISTS/IN subquery comparison expression.

CorrelatedColumnInfo

Stores specific metadata for columns that cross the correlation boundary.

pub struct CorrelatedColumnInfo {
    pub col: Column,
    pub field: FieldRef,
    pub depth: usize,
    pub delim_scan_node_id: usize, // References the domain scan provider map (for compiling DelimGets)
}

DependentJoinRewriter

The stateful tree rewriter that implements TreeNodeRewriter to execute Phase 1.

pub struct DependentJoinRewriter {
    current_id: usize,                             // Unique ID assigned to each visited node in the logical plan
    subquery_depth: usize,                         // Tracks active subquery nesting depth
    nodes: IndexMap<usize, Node>,                  // Temporary metadata tracking for each visited LogicalPlan
    stack: Vec<usize>,                             // Traversal path (node IDs) from root to current node
    unresolved_outer_ref_columns: IndexMap<Column, Vec<ColumnAccess>>, // Holds column accesses currently unlinked to providers
    alias_generator: Arc<AliasGenerator>,          // Generates unique aliases for rewritten subqueries
    pub domain_columns_provider_nodes: IndexMap<usize, LogicalPlan>, // Maps node ID to its corresponding column provider node
}

ColumnAccess

Tracks an active column access before it is resolved by its provider.

struct ColumnAccess {
    stack: Vec<usize>,             // Accessor's traversal path to root
    node_id: usize,                // Node ID of the accessing operator
    col: Column,                   // Column being accessed
    field: FieldRef,               // Column type/field info
    subquery_depth: usize,         // Correlation depth of the access
    provider_node_id: usize,       // Node ID of the provider operator (resolved at LCA)
}

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).

  1. Phase 2 (DependentJoinDecorrelator): We will introduce a top-down decorrelation algorithm to eliminate DependentJoin nodes, turning them into standard physical joins (e.g., Left, Semi, Anti Joins).
  2. Phase 3 (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.

@github-actions github-actions Bot added sql SQL Planner logical-expr Logical plan and expressions optimizer Optimizer rules core Core DataFusion crate substrait Changes to the substrait crate common Related to common crate labels Apr 11, 2026
@duongcongtoai duongcongtoai changed the title Feat dependent join integration Introduce dependent join LogicalPlan to support complex subquery decorrelation Apr 11, 2026
@alamb

alamb commented Apr 27, 2026

Copy link
Copy Markdown
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
duongcongtoai marked this pull request as ready for review May 4, 2026 13:10
@github-actions

Copy link
Copy Markdown

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.

@github-actions github-actions Bot added the Stale PR has not had any activity for some time label Jul 24, 2026
@github-actions github-actions Bot closed this Jul 31, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

common Related to common crate core Core DataFusion crate logical-expr Logical plan and expressions optimizer Optimizer rules proto Related to proto crate sql SQL Planner Stale PR has not had any activity for some time substrait Changes to the substrait crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants