Skip to content

Proposal: Unified Subquery Deduplication and Merging #1085

Description

@kagamiori

1. Problem

Deduplication

When a query contains structurally identical subqueries, each is translated and
decorrelated independently, producing redundant joins that scan and aggregate
the same data multiple times.

SELECT * FROM nation n
WHERE n_nationkey > (SELECT avg(n_nationkey) FROM nation s
                     WHERE s.n_regionkey = n.n_regionkey)
  AND n_regionkey < (SELECT avg(n_nationkey) FROM nation s
                     WHERE s.n_regionkey = n.n_regionkey)

Both subqueries compute avg(n_nationkey) GROUP BY n_regionkey over the same
table. Without dedup, this produces two LEFT JOINs. After dedup, one suffices —
both filter predicates reference the same result column.

Merging

When subqueries share the same FROM table, same correlation, and same GROUP BY
but differ only in their aggregates, each produces a separate join that scans
and groups the same data.

SELECT * FROM nation n
WHERE n_nationkey > (SELECT avg(n_nationkey) FROM nation s
                     WHERE s.n_regionkey = n.n_regionkey)
  AND n_regionkey < (SELECT max(n_nationkey) FROM nation s
                     WHERE s.n_regionkey = n.n_regionkey)

Without merge, this produces two LEFT JOINs with separate aggregations. After
merge, one LEFT JOIN with a combined aggregation (avg + max) suffices.

Dimensions of applicability

The following dimensions determine whether two subqueries can be deduplicated or
merged:

By subquery type:

Type Dedup eligible? Merge eligible? Reason
Scalar Yes Yes Result is a single value — duplicates produce identical values; compatible aggregates can be combined.
IN Yes No The projected column is the semi-join equality key. Different projections mean different join conditions — merging would change semantics (requiring both conditions to match the same row).
EXISTS Yes No The projection is trivial (SELECT 1). Identical FROM + WHERE is already dedup. Different WHERE conditions are semantically different predicates, not mergeable.

By correlation:

Correlation Dedup eligible? Merge eligible? Reason
Uncorrelated Yes Yes No outer dependency. Same computation → same result.
Correlated Yes, when same outer references Yes, when same outer references A correlated subquery is a pure function of the correlation key value and the database state. Same outer references → same function → dedup/merge is correct regardless of outer-table pre-filtering.

By clause location:

Location Dedup eligible? Merge eligible? Reason
Same expression (e.g., WHERE a > sub1 AND b < sub2) Yes Yes Both subqueries are in the same expression tree — visible together.
Different expressions in same node (e.g., separate projection columns) Yes Yes Same outer table, same scope. Visible together if expressions are batched.
Across filter and projection (same query block) Yes Yes Same outer table, same correlation key values. The subquery result is a pure function independent of which clause references it.
Across UNION ALL branches Not safely Not safely Each branch has its own outer table instance. The decorrelated join edge is bound to a specific outer table in a specific DT. Reusing a result column from one branch's DT in another creates dangling references. This is a structural constraint of the query graph model, not a semantic issue.

2. How Other Engines Handle Dedup and Merge

Overview

Feature Presto (0.297) DuckDB (1.4.4) Spark (4.0) Axiom (proposed)
Dedup Syntactic identity No Canonicalized plan Positional structural comparison
Merge No No Uncorrelated scalar only Correlated + uncorrelated scalar
When Planning time Post-decorrelation Pre-decorrelation (pre-pass)
Handles alias differences No Yes (canonicalized) Yes (positional, ignores names)

Presto

Deduplicates subqueries via syntactic identity at planning time. In
SubqueryPlanner.java, before creating a new join node (LateralJoinNode for
scalar subqueries, ApplyNode for IN/EXISTS), checks if the expression already
has a mapping in the TranslationMap. Two subqueries that
are semantically identical but syntactically different (e.g., different column
aliases) are NOT deduplicated. No merging support — each subquery becomes a
separate LateralJoinNode and is decorrelated independently.

DuckDB

No deduplication. BoundSubqueryExpression::Equals() always returns false,
preventing the CommonSubExpressionOptimizer from detecting identical
subqueries. No merging support.

Spark

Deduplicates via canonicalized plan comparison in MergeSubplans (formerly
MergeScalarSubqueries). Canonicalization normalizes column names, so
semantically equivalent subqueries with different aliases are deduplicated.
Supports merging uncorrelated scalar subqueries with the same FROM but different
aggregates — tryMergePlans recursively matches plan structures and combines
aggregate lists. Merging is restricted to uncorrelated scalars because the rule
runs after decorrelation — correlated subqueries have already been rewritten
into LEFT JOINs and are no longer visible as subquery expressions. Merged
results are wrapped in a CTE with GetStructField extraction.

Axiom (proposed)

Operates before decorrelation (pre-pass over the logical plan tree), so it
naturally sees both correlated and uncorrelated subqueries. Uses positional
structural comparison (subqueryExprEquivalent) that resolves column references
by ordinal position, making it robust to planner-generated name suffixes. This
is the only engine that supports merging correlated scalar subqueries.


3. Current Architecture: How Subqueries Are Processed

ToGraph::makeQueryGraph() walks the logical plan tree top-down and builds a
flat query graph (DerivedTable with tables, join edges, and conjuncts). When
it encounters a Filter or Project node, it calls processSubqueries to extract
and decorrelate subqueries.

Processing flow per call:

  1. extractSubqueries(expr) — walks the expression tree, classifies each
    subquery into scalars, inPredicates, or exists.
  2. Optionally wraps currentDt_ via finalizeDt (to prevent self-referencing
    join edges when prior decorrelated joins exist).
  3. For each subquery: translateSubquery creates a new DT, recursively calls
    makeQueryGraph, then the correlated/uncorrelated handler creates the
    appropriate join edge.
  4. Stores the result in subqueries_[exprPtr] = resultColumn, used later by
    translateExpr to replace subquery expressions.

Key constraint: processSubqueries is called per-expression (one for each
projection column, one for the filter predicate). Subqueries in different
expressions are in different calls, separated by potential finalizeDt DT
wrapping. This makes cross-call optimization non-trivial — column references
become stale after wrapping and must be rewritten via exportExpr.


4. Proposed Approach

(A proof-of-concept prototype is available in #1187)

Core idea

Move all deduplication and merging to a single pre-pass over the logical
plan tree, before makeQueryGraph. The pre-pass sees all subqueries globally
and builds a map (subqueryMap_) that records which subqueries are duplicates
and which are mergeable. During query graph construction, processSubqueries
applies this pre-computed map via pointer-identity lookups.

Architecture

┌─────────────────────────────────────────────┐
│  buildSubqueryMap (pre-pass)                │
│                                             │
│  1. collectAllSubqueries — walk plan tree,  │
│     gather all SubqueryExpr/IN/EXISTS       │
│                                             │
│  2. Dedup — O(n²) comparison via            │
│     subqueryExprEquivalent (all types)      │
│                                             │
│  3. Merge — group remaining scalars by      │
│     isSameMergeKey, build merged plans      │
│                                             │
│  Output: subqueryMap_                       │
│    duplicate → {representative, 0}          │
│    merge original → {merged repr, index}    │
└─────────────────────────────────────────────┘
                    │
                    ▼
┌─────────────────────────────────────────────┐
│  processSubqueries (per-call processing)    │
│                                             │
│  1. Extract subqueries from expression      │
│                                             │
│  2. Cross-call check: if exprPtr already    │
│     in subqueries_, skip (rewrite via       │
│     exportExpr if finalizeDt wrapped)       │
│                                             │
│  3. Apply subqueryMap_: replace originals   │
│     with representatives                    │
│                                             │
│  4. finalizeDt if needed                    │
│                                             │
│  5. Process remaining subqueries            │
│     (translateSubquery + decorrelation)     │
│                                             │
│  6. Store results for all mapped originals  │
│     in subqueries_                          │
└─────────────────────────────────────────────┘

Key data structures

Data structure Type Role
subqueryMap_ F14FastMap<ExprPtr, SubqueryMapEntry> Pre-computed map from each duplicate/mergeable original to its representative and output index. Populated once by buildSubqueryMap.
SubqueryMapEntry {ExprPtr representative, size_t outputIndex} For dedup: representative is an original ExprPtr, outputIndex is 0. For merge: representative is a synthetic SubqueryExprPtr wrapping a merged AggregateNode, outputIndex identifies the aggregate column.
subqueries_ F14FastMap<ExprPtr, ExprCP> Maps subquery expressions to their replacement columns. Populated during processing. Used by translateExpr and by cross-call pointer-identity dedup.

Comparison functions

Function Purpose Scope
subqueryExprEquivalent Unified dedup comparison for all types. Dispatches internally: scalar → isSameComputation on plans; IN → compares form + left key + subquery plan; EXISTS → compares form + subquery plan. Dedup
isSameMergeKey Compares two scalar subquery plans for merge compatibility: identical child structure and grouping keys, but aggregate functions may differ. Merge
isSameComputation Recursively compares LogicalPlanNode trees positionally (types by position, column references by ordinal). The foundation both above functions build on. Both

How cross-call dedup/merge works

When the pre-pass maps subqueries A (in filter) and B (in projection) to the
same representative R:

  1. Filter's processSubqueries: extracts A → subqueryMap_ lookup finds
    R → processes R → stores subqueries_[A] and subqueries_[B] (all
    originals mapped to R get their results stored immediately).

  2. finalizeDt wraps currentDt_ (because a non-inner join was added).

  3. Projection's processSubqueries: extracts B →
    subqueries_.contains(B) is true (stored in step 1) → B is added to
    alreadyProcessed → its stale column reference is rewritten via
    exportExpr through the wrapper. No new join created.

This reuses the existing exportExpr rewriting infrastructure that was
already in place for pointer-identity cross-call dedup.

Merged plan construction

For a merge group of scalar subqueries with the same merge key:

  1. buildMergedPlan creates a new AggregateNode that combines all group
    members' aggregates, reusing the representative's child subtree.

  2. Aggregate expressions from non-representative members have their
    InputReferenceExpr names remapped to the representative child's column
    namespace via remapAggregateInputNames.

  3. The merged node is registered in the subfield tracker via
    registerAllChannelsAsUsed (since it's a synthetic node not seen by the
    original subfield tracking pass).

Scope boundaries

collectAllSubqueries stops at set operation nodes (kSet — UNION, INTERSECT,
etc.). buildSubqueryMap recurses into each branch independently. This ensures
subqueries from different union branches are never grouped together — each
branch is an independent scope with its own currentDt_ during
makeQueryGraph.


5. Capability Summary

What can be deduplicated

Dimension Supported? Reason
Scalar subqueries Yes Full structural comparison via isSameComputation.
IN subqueries Yes Compares form + left key + subquery plan via subqueryExprEquivalent.
EXISTS subqueries Yes Compares form + subquery plan via subqueryExprEquivalent.
Correlated Yes Outer references are part of the plan tree; structural comparison covers them.
Uncorrelated Yes No outer dependency — same computation → same result.
Same expression Yes Extracted into the same batch.
Across projection columns Yes Projection expressions batched into a single processSubqueries call.
Across filter and projection Yes Pre-pass sees both globally; cross-call exportExpr rewriting handles finalizeDt wrapping.
Across UNION branches No Each branch has its own outer table instance; join edges are DT-specific. Reusing columns across branches creates dangling references. This is a structural constraint of the query graph model.

What can be merged

Dimension Supported? Reason
Scalar subqueries Yes Aggregates from different subqueries combined into one AggregateNode.
IN subqueries No (semantic) The projected column is the semi-join equality key. Different projections → different join conditions → different semantics.
EXISTS subqueries No (semantic) Projection is trivial (SELECT 1). Identical FROM + WHERE → dedup. Different WHERE → semantically different predicates.
Correlated Yes Pre-pass operates before decorrelation, so correlated subqueries are still visible as SubqueryExpr nodes.
Uncorrelated Yes Same as correlated, but uncorrelated global aggregations over bare scans are excluded to preserve constant folding (merging would replace two independently-foldable literals with a cross-join).
Same expression Yes Pre-pass sees all subqueries globally.
Across projection columns Yes Same as above.
Across filter and projection Yes Same as above — pre-pass + cross-call exportExpr rewriting.
Across UNION branches No Same structural constraint as dedup.

Comparison with other engines

Capability Presto DuckDB Spark Axiom (proposed)
Dedup scalar Syntactic only No Uncorrelated only All (correlated + uncorrelated)
Dedup IN Syntactic only No No Yes
Dedup EXISTS Syntactic only No No Yes
Dedup across clauses No No No Yes (filter + projection)
Merge scalar No No Uncorrelated only All (correlated + uncorrelated)
Merge across clauses No No No Yes (filter + projection)
Handles alias differences No No Yes Yes

Metadata

Metadata

Assignees

Labels

No labels
No labels

Type

No type

Projects

No projects

Milestone

No milestone

Relationships

None yet

Development

No branches or pull requests

Issue actions