Repository navigation
feat(ir): add logical sub-DAG sharing and export without execution timing - #537
Conversation
a59aff8 to
8f384f7
Compare
9369217 to
9870df6
Compare
8f384f7 to
d8f0b1a
Compare
9870df6 to
88c2931
Compare
d8f0b1a to
dd71af0
Compare
c0638f6 to
d414efd
Compare
dd71af0 to
7d5c382
Compare
d414efd to
ec0417b
Compare
7d5c382 to
6aba8f2
Compare
00a3d27 to
467baad
Compare
467baad to
aecd114
Compare
ba26c8b to
7cd4ee0
Compare
76b23bc to
c30977a
Compare
60a9e4f to
cb00197
Compare
cb00197 to
700838b
Compare
cbc4cf0 to
0f6f539
Compare
I cannot understand the first few paragraphs of this PR According to my understanding. The following sounds more natural: |
Selvomega
left a comment
There was a problem hiding this comment.
CSE and canonicalize roughly lgtm (though I also feel some of their code snippet quite sketchy)
I feel something is off about the export part. My feeling is
- At the first place I don't know why we need exporting logical DAG here, since we have decided that it is not the final output.
- To make the LogicalDAG exportable, I also don't think duplicating types (as what we are doing for now in
crates/types/src/ir/wire.rs) is right. There should be more light-weight method for it.
0f6f539 to
fa9c414
Compare
700838b to
f2b7b3f
Compare
Thanks for pointing out the pr body unclear writing, I have updated it. |
Follow #537: the remaining tests read Operator<NodeId> payloads from the physical DAG, and docs name ir::flat / ir::physical_export instead of the removed ir::export. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
Thanks for the review. Addressed in 4ac03c3:
|
Follow #537 and #541: the physical DAG carries Operator<NodeId> instead of the removed wire mirror types. Port the moved physical_planner, the planner's lifecycle export and the tests that read or hand-build physical DAGs; hand-built nodes now name their children by id, matching their edges. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
4ac03c3 to
ebea327
Compare
f2b7b3f to
3f5d369
Compare
Follow #537: the remaining tests read Operator<NodeId> payloads from the physical DAG, and docs name ir::flat / ir::physical_export instead of the removed ir::export. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Follow #537: the remaining tests read Operator<NodeId> payloads from the physical DAG, and docs name ir::flat / ir::physical_export instead of the removed ir::export. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
3f5d369 to
52b4003
Compare
ebea327 to
a419b14
Compare
|
This PR shows as merged, but only into its old stack base |
Follow #537: the remaining tests read Operator<NodeId> payloads from the physical DAG, and docs name ir::flat / ir::physical_export instead of the removed ir::export. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Follow #537: the remaining tests read Operator<NodeId> payloads from the physical DAG, and docs name ir::flat / ir::physical_export instead of the removed ir::export. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Rebased on main d4869a7 (DF 54). Revised after review: the logical export and the
wire.rsmirror types are removed (see item 3).Problem
#536 defined the unified
ir::OperatorNodeDAG, and #567 and #560 added summary coverage andSummaryMergeto it. But canonicalization, common sub-DAG sharing (CSE) and a flat, serializable form have not been implemented on these new types. They exist only for the old split IRs:pre_asap::canonicalize::canonicalize(QueryExpr) -> QueryExprQueryExprpre_asap::cse::share_common_sub_dags(Vec<(Id, QueryExpr)>)QueryExprpost_asap::cseSummaryNodedag_export::export/export_post_asapQueryExpr(+ post-ASAP substitutions)So, on
OperatorNode:Two equivalent spellings stay different shapes. Example:
Filter { EXISTS (s) }andJoin { Semi, true, left, s }mean the same thing but cannot be matched structurally.Two identical sub-DAGs stay two
Rcs. Example: one query with the same grouped p50 on both sides of a comparison:Nothing collapses A and B into one node with two consumers.
A shared DAG cannot be written out flat.
OperatorNodeserializes as a tree, so every shared sub-DAG would be repeated. The earlier version of this PR fixed that with a second copy of the operator and scalar types (wire.rs:NonASAPOpKind,WireScalarExpr, ...) whose children were ids.Therefore, this PR adds:
ir::canonicalize: rewrites equivalent spellings into one shape (fixes 1).ir::cse: shares structurally identical sub-DAGs across a workload batch (fixes 2). This is the identical-expression rule of docs: propose workload-wide planning, summary sharing, and materialization #509 Pass 2: "The input and computation semantics are identical" → "One common computation node serving multiple consumers."Operator,NonASAPOp,ASAPOp,ScalarExpr,Predicate,ProjectItem,SortKeyandQueryRoottake the child reference as a type parameter, defaulting toRc<OperatorNode>.ir::flat::flattenwrites a DAG as a list of nodes whose operators areOperator<NodeId>. No logical export document is added: the planner's stages passRc<OperatorNode>DAGs directly, and the physical DAG (feat(runtime): compile unified operator and scalar graphs #541) is the first consumer of the flat form.It does not implement the summary-capability rule or the window-composition rule of #509 Pass 2, Pass 1 candidate generation, or any physical-stage decision (materialization, timing assignment, phase splitting).
Proposed method
All on bound
OperatorNodeDAGs incrates/types/src/ir/.1. Canonicalize (
ir::canonicalize). Bottom-up over every child (operator inputs and nodes read by scalar expressions), memoized byRcpointer so a shared sub-DAG is rewritten once. At each node it applies local rules until none matches:Limit { n: k, offset: 0, Sort { one descending column key, [passthrough Project] Aggregate(Reduce(by), [Count | Sum]) } }becomesAggregate(Reduce(partition), [TopK { k }])over the unchanged inner aggregate. It is skipped for ascending sorts, offsets, ranking by a group key, a Limit partition that differs from the Sort partition, andSumover a derived childAggregate(rate/increase).Filter, oneEXISTS (s)conjunct becomesJoin { Semi, true },NOT EXISTS (s)becomesJoin { Anti, true }, and a positivex IN (s)becomesJoin { Semi, x = Column(left_width) }. Remaining conjuncts stay in an outerFilter. One conjunct per round.NOT IN, anINwhose probe contains a scalar subquery, scalar subqueries (a cross join would lose their zero-row NULL and multi-row error semantics), andROW_NUMBER()filters (removing the window column would change the visible schema).Untouched sub-DAGs keep their pointer identity, a rewritten shared sub-DAG stays shared, and the pass is idempotent. A
Concatwhose first branch changes schema drops itsdiscriminator_unique_key.2. Share common sub-DAGs (
ir::cse). Hash-consing over all roots of a workload batch, bottom-up:PromqlScalarFromVector,ScalarSubquery,Exists,InSubquery). Soscalar(v)in two queries can sharev.structural_hashonly picks a bucket. The decision is typed equality (same_node): children by pointer (or memoized value comparison), thenresult_kind,schema,timing,coverage,guarantee, and the operator's own fields (operator.map_children(|_| ()), so children never enter the comparison twice). Guarantee and own fields must also serialize identically, so0.0vs-0.0andNaNvsNaNare not shared.schema.has_unique_key(). ASAP nodes have no such gate.coverageis part of both hash and equality, so two summary states over different observations are never shared.PartialEq-equal to its input); only itsRcstructure may now alias.The pass does not decide whether sharing is cheaper. That is selection's job (#509 stage 3).
3. Generic child references and
ir::flat. Each operator and scalar type gets a parameterCfor its child reference, defaulting toRc<OperatorNode>, so existing code is unchanged. The traversal helpers (children,scalar_exprs,operator_refs,kind_name) work for anyC;map_children/map_operator_refsmapCto a new typeD. Schema derivation and validation stay on the defaultRc<OperatorNode>form.flatten(roots)visits the DAG children-first and gives each distinctRc(by pointer) the next index. Each node becomes aFlatNode: the same fields asOperatorNode, withoperator: Operator<NodeId>. A scalar root stays aQueryRoot::Scalar(ScalarExpr<NodeId>).OperatorNodeitself is not generic: a self-referential default (OperatorNode<C = Rc<OperatorNode>>) is rejected by rustc (E0391).Key code interfaces
crates/types/src/ir/node.rs,non_asap.rs,asap.rs,scalar.rs,query.rscrates/types/src/ir/flat.rscrates/types/src/ir/canonicalize.rs,cse.rsUsage:
Examples
Example 1: a shared node is flattened once (
crates/types/src/ir/flat.rstests)Input: two roots
Filter(true), each over the sameValuesRc.flattenreturns 3 nodes:0 = Values,1 = Filter { child: 0 },2 = Filter { child: 0 }, and roots[Operator(1), Operator(2)](a_shared_node_is_flattened_once). AProjectwhosechildandScalarSubqueryare the sameValuesgives 2 nodes, withchild: 0andcols[0].expr == ScalarSubquery(0)(scalar_references_become_node_ids). A constant scalar root gives 0 nodes (a_constant_scalar_root_has_no_nodes), and the JSON round-trips with everytiming = None(json_round_trips).Example 2: CSE on the identical-expression rule (
crates/types/src/ir/cse.rstests)The
BinaryOpfrom the Problem section, aftershare_common_sub_dags(vec![("q", root)]):lhsandrhsare the sameRc(single_query_shares_its_own_repeated_sub_dag). The grouped aggregate has unique keyby [1], so it is legal to share.Quantile q=0.5 col=2in two roots (median_and_explicit_half_percentile_merge)col=2vscol=3(distinct_column_quantiles_do_not_merge)no_unique_keys_means_no_merge_even_when_structurally_identical)without(..)aggregate (group_keys_gate_still_prevented_when_partition_by_without_used)Dedup { cols: [1] }(dedup_gates_sharing_the_same_as_aggregate)Dedupadds a unique keyvector(scalar(sum by (service)(up)))twice (scalar_referenced_vector_is_shared_across_queries)SummaryAgg,alpha=0.01, same guarantee (asap_nodes_share_without_a_unique_key)alpha=0.01vs0.001; exact guarantee vsNone(asap_nodes_with_distinct_parameters_or_guarantees_are_not_shared)SummaryEstimatep95 and p99 over equal producers (evaluations_share_their_producer_but_not_each_other)q0.0vs-0.0,+infvs-inf,NaNvsNaN(signed_zero_and_nonfinite_values_remain_distinct)+infvs+infis sharedcoveragecoverageis in hash and equality (no dedicated test)Example 3: canonicalization (
crates/types/src/ir/canonicalize.rstests)Limit 5 { Sort desc col 1 { Aggregate(by [1], [Count]) } }Aggregate(by [], [TopK{k:5}]) { Aggregate(by [1], [Count]) }(promotes_count_ranked_limit_sort); also through a passthroughProjectFilter { EXISTS (s) }Join { Semi, true, left, s }(exists_filter_becomes_semi_join)Filter { NOT EXISTS (s) }Join { Anti, true, left, s }Filter { service IN (s) }over a 3-column leftJoin { Semi, service = Column(3) }Filter { value > 1 AND EXISTS (s) }Filter { value > 1 } { Join { Semi } }Filter { NOT IN (s) }Project [ScalarSubquery(s) AS a, ScalarSubquery(s) AS b]s(scalar_subqueries_remain_explicit_and_shared)Filter { rn <= 5 } { ROW_NUMBER window }Every rewrite test also checks
canonicalize(out)returns the sameRc.Out of scope
flatten.MaintainPopulation.populationstill holds its input as an inlineRc<OperatorNode>, so the flat form embeds that subgraph rather than linking it by id.pre_asap::canonicalize,pre_asap::cse,post_asap::cseanddag_export.Stack and validation
Revised logical foundation 3/6 · Previous: #560 · Next: #539 · Tracker: #528
Order: #567 → #560 → #537 → #539 → #540 → #541 → #542 → #543.
flatten: its payload isOperator<NodeId>, the runtime's 374-line conversion from the wire types became amap_childrencall, andEdgeRole/GroupingEdgeCompatibilitymoved there. feat(runtime): compile unified operator and scalar graphs #541–refactor(ir): remove legacy scaffolding and finish tooling migration #543 each carry one port commit.🤖 Generated with Claude Code