Repository navigation
Add AQE to DataFusion #23194
Description
Activity
Note there is a related ticket / work in balllista
Similar work for distributed-dafafusion from @gabotechs: datafusion-contrib/datafusion-distributed#486
I only looked at the issue/PR description for this issue, but I can already see many commonalities:
- pipeline breaker/boundaries as a good place to accumulate runtime statistics
- re-use of the built-in statistics propagation mechanism (great for re-use), only fueled with runtime statistics
- runtime statistics must be sampled as we are in a streaming computational model (idea behind
SamplerExecin the above PR, and the similar buffer node here)
I wonder how much can be re-used across core DF/distributed DF/ballista, there are different challenges and the same logical concept has different forms in the three cases, but the mechanism seems very similar, if not identical.
WDYT?
Reacted by Brent Gardner and Andy GroveHow about starting in datafusion-contrib/datafusion-aqe?
I think this can be implemented almost entirely outside of the core crates, where we are already struggling with a large umber of features.
In
datafusion-distributed, once this EPIC is finished:And once the worker protocol has been fully abstracted out of gRPC so that there can be an in-memory implementation:
- feat: added a pluggable
WorkerTransport, with Arrow Flight optional. datafusion-contrib/datafusion-distributed#508 (driven by ParadeDB folks)
There should be a fully usable single-node AQE implementation that works pretty much out of the box in non-distributed contexts, without requiring materialization to disk. It still requires all the coordination machinery from
datafusion-distributedthough, so I don't think it can be easily factored out.Reacted by Andrew Lamb and Edmondo Porcu- feat: added a pluggable
There should be a fully usable single-node AQE implementation that works pretty much out of the box in non-distributed contexts,
FYI @milenkovicm and @avantgardnerio -- see @gabotechs comments above
Just for context, https://www.cs.cmu.edu/~15721-f24/papers/AQP_in_Lakehouse.pdf
I'm not sure we can make one common AQE but we should definitely try. If we could make it as pluggable and flexible as planning rules are maybe it will work.
Reacted by Andrew Lamb and Brent GardnerOne thing that I'd expect to be shared across all approaches and probably very key to all of them, is to ensure we can reuse the current
PhysicalOptimizerRuleimplementations that swap joins, etc... for the during-execution passes that adapt the query at runtime.I imagine this implies that these
PhysicalOptimizerRules should be idempotent among other things.Reacted by Brent Gardner and Andrew LambI'm not sure we can make one common AQE but we should definitely try.
I'd like to highlight that Coralogix runs Ballista internally for distribution, so this PR was definitely written with distribution in mind.
The key architectural trait I'd like to highlight is the original "isomorphic scaling" approach that Datafusion & Ballista took from inception: the same mechanisms Datafusion uses to spread load across cores is the same mechanism that Ballista uses to distribute across executors.
The intent here was to follow that pattern: Stage boundaries inserted by these rules would become shuffle writes in Ballista. (The additional parallel window function PRs all follow this pattern as well).
So I don't think this is a competing approach, and a
datafusion-aqecrate would just end up being a 3rd effort that would diverge rather than converge existing approaches. What this PR hopes to do is add the necessary primitives within Datafusion to make distributed AQE easier.I'm not as familiar with
datafusion-distributed, but I'm going to spend the rest of my day comparing the other two implementations and roadmaps with what is provided here - the goal being to find commonality between them and move it in-repo.Reacted by Andrew LambPhysicalOptimizerRule
The
RuntimeRuletrait only exists in the PoC because of a circular dependency issue that is easily resolved if we agree on the direction. The intent is to usePhysicalOptimizerRule- though each one would have to be made "runtime aware" and added to theRuntimeOptimizerwhite list of rules; but, this approach can happen incrementally by converting one rule at a time, not a "big bang" where we deliver no value until every impl changes.I'm on a vacation will help once I'm back
Reacted by Brent Gardner, Nathaniel Cook and Andrew LambThe key architectural trait I'd like to highlight is the original "isomorphic scaling"
I found this terminology somewhat complicated to grok at first -- in my words this means that the work is divided across DataFusion partitions and additional parallelization is implemented explicitly in terms of partitions (rather than more cores for a particular operator, for example).
The rationale is that ballista and datafusion distributed reason about and can distirbuted partitions across cores and nodes
@gabotechs , I've reviewed dfd's architecture and your above proposal, and I think we can come to a common proposal that satisfies DF, dfd, and Ballista.
Ideas:
- The PoC's StageBoundary becomes the same shape as SamplerExec, parameterized by a single capacity N. Small N = sampler-shaped (bounded peek with NDV/null-pct/velocity); N=∞ = pipeline breaker (exact stats, gated on spill). Same operator, one knob: streaming-vs-eager collapses into a scalar (with other possible tricks down the road like memory-pressure awareness, spill, and upstream pipeline breaker awareness)
- I think
SamplerExecprovides truly richer stats, I think we can build stats that are a superset ofLoadInfoand thePartitionStatisticsof today (adding PartitionSortExtrema while we are at it) - The StageBoundary/SamplerExec could take a listener: Option<Arc>, in the constructor. DF would use None, but upstream projects could get async notifications
This way dfd's coordination layer stays completely intact — its oneshot→gRPC delivery becomes the StatsListener impl — and DF picks up the observation primitive natively. Do you think this would be a fruitful direction? If you'd like to merge in SamplerExec then DF would have the perfect slot for it. Or if you prefer, the work would be lifted directly with attribution.
As an aside, I would like to correct a confusing sentence from before though - when I said:
What this PR hopes to do is add the necessary primitives within Datafusion
I meant primitives + intra-DF-AQ-execution. This way multiple projects don't need to redo the work of AQE aware optimizations, the entire DF community would help in the AQE project. And what were previously single-core operations (windows without partitionby) within DF become parallelized within a worker/executor.
WDYT?
I forgot about your comment, then went and re-invented most of it in the reply above ☝️ . I'm glad we independently arrived at similar conclusions.
Reacted by Alessandro SolimandoAdaptive aggregation strategy — pick between hash and sort aggregation, or between Partial and Single, based on observed cardinality.
Related PR in progress , and we have good benchmark result.
Reacted by Brent Gardner and Daniël HeresThanks @zhuqi-lucas !
I'm trying to compile a list of the ongoing AQE-like efforts, inside the existing DF framework, including my own:
adaptive agg(feat(aggregate): cost-aware partial-aggregation skip (opt-in) #22518) - sampling stats, makes runtime decisions, within operatorparallel sort(feat: parallel sort-preserving merge (PSRS) for >2x merge speedup (demo) #23124) - sampling stats, makes runtime decisions, within operatorhalo rows(Parallel bounded RANGE-frame window functions without PARTITION BY (draft) #23026) - full materialization stats from existing operator, runtime decisions, between SortExec & BWAGprefix scan([PoC] Parallel prefix scan for cumulative-aggregate window functions coralogix/arrow-datafusion#426) - full mat stats, runtime, between SortExec & BWAG
None of these have been merged.
adaptive agghas positive feedback regarding moving more decisions into the operator.halo rowshas negative feedback about moving information between operators (the inspiration for this issue).- another drop BWAG thread also mentions moving more concerns (spilling) into operators
I think either architectural direction would be fair:
- move all AQE related efforts into operators
- extract an AQE specific framework like described in this issue.
I don't understand a principled 3rd option, and I'd prefer we didn't implicitly set one without the explicit understanding that the issue has been raised - anything we do henceforth is conscious.
I don't have a strong opinion one way or the other (my goal is parallel windows), but the journey has led me to believe that AQE is foundational in its effects on optimizer rules and operators (you either end up with separate static and runtime, or hybrid that does both). This is true within datafusion and in downstream projects (ie. duplication will, is, and has happened )
CC @alamb @andygrove
12 remaining items
Thanks @avantgardnerio, what is the scope for this AQE? Currently discussing datafusion and distributed variants.
I have some feeling that support AQE for datafusion single machine engine and distributed variants are totally different amount of work.
For distributed the stats synchronization is supposed to be more complicated, it would require some coordinator that knows about all hash sizes, allocations, stats, etc.
Also to insert pipeline breakers the distributed prob should support some sort of stages like Spark.
Is the scope to come up with some universal trait/approach that would shape the AQE for single machine and distributed?
@comphead if you click through into the (working) reference PoC, you'll find numbered stages:
01)RuntimeOptimizerExec 02)--HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(id@0, group_key@0)], projection=[group_key@2, sum_payload@3, payload@1] 03)----CoalescePartitionsExec 04)------DataSourceExec: partitions=4, partition_sizes=[1, 0, 0, 0] 05)----ProjectionExec: expr=[group_key@0 as group_key, sum(big.payload)@1 as sum_payload] 06)------PipelineBreakerBuffer 07)--------AggregateExec: mode=FinalPartitioned, gby=[group_key@0 as group_key], aggr=[sum(big.payload)] 08)----------RepartitionExec: partitioning=Hash([group_key@0], 4), input_partitions=4 09)------------PipelineBreakerBuffer 10)--------------AggregateExec: mode=Partial, gby=[group_key@0 as group_key], aggr=[sum(big.payload)] 11)----------------DataSourceExec: partitions=4, partition_sizes=[4, 3, 3, 3]They just execute in an existing
Arc<dyn PhysicalPlanNode>, being re-written at boundaries (obviously not the in-progress ones).some coordinator that knows about all hash sizes, allocations, stats
See
RuntimeOptimizerExecFor distributed the stats synchronization
For Ballista, I assumed this would just be 1. recognize the stage boundary, 2. send the info to the scheduler on the stage completion gRPC call (admittedly this breaks down with streaming or large responses). I was hoping we could start with a dedicated stats collector node outside any operator (
PipelineBreakerBufferdoes this with row_count in the PR,SortExtremain another, no reason not to add more - KLL, etc), but it's an interesting question if we can crack open existing operators and extract their state, or if these concerns are intractable.Is the scope to come up with some universal trait/approach that would shape the AQE for single machine and distributed?
That is what I interpreted I was being asked to do by Andrew, yes. This adds true AQE to DF without much disruption (comparatively speaking :), and defers the distribution concerns (and distribution only) to downstream implementations.
Edit, sorry, numbered stages:
01)RuntimeOptimizerExec 02)--HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(id@0, group_key@0)], projection=[group_key@2, sum_payload@3, payload@1] 03)----StageBoundaryBuffer: stage=0 04)------CoalescePartitionsExec 05)--------DataSourceExec: partitions=4, partition_sizes=[1, 0, 0, 0] 06)----StageBoundaryBuffer: stage=0 07)------ProjectionExec: expr=[group_key@0 as group_key, sum(big.payload)@1 as sum_payload] 08)--------AggregateExec: mode=FinalPartitioned, gby=[group_key@0 as group_key], aggr=[sum(big.payload)] 09)----------RepartitionExec: partitioning=Hash([group_key@0], 4), input_partitions=4 10)------------AggregateExec: mode=Partial, gby=[group_key@0 as group_key], aggr=[sum(big.payload)] 11)--------------DataSourceExec: partitions=4, partition_sizes=[4, 3, 3, 3]Reacted by Oleks VWe (Datadog), will gladly donate https://github.com/datafusion-contrib/datafusion-distributed to Apache if that implies hosting it as a new crate in https://github.com/apache/datafusion. We've built the crate leaving that door open from the beginning, both from a philosophy and code standpoint.
We've discussed this in the past, but at that moment it was not the right time. Now that we've been running it in production at a huge scale for a while and we can no longer break things and move fast, the door is open.
I am much more in favor of adopting datafusion-distributed compared to starting with a new crate in the datafusion repo --
datafusion-distributedhas had a lot of thought in it and I have heard from many users they are using it / plan to use it (as are we at InfluxData). This gives me confidence that its APIs are broadly applicable and we can maintain it.To be clear, datafusion-distributed is not a complete distributed query engine, it is more like "the common pieces needed to build one" -- like shuffles and network communication traits.
I'm going to study datafusion-distributed over the next couple of days to understand the architecture and how it varies with Ballista so that I can make sure that I am not proposing something that would help with Ballista and not with datafusion-distributed.
If we could somehow migrate Ballista to use datafusion-distributed I think that would also be amazing (and would remove all doubt about moving datafusion-distributed into the core datafusion crate)
Reacted by Gabriel, Oleks V and Daniël HeresOne point I want to make is that dynamic optimization of different query stages (a.k.a. AQE in Spark but used in some database engines as well in a more ad-hoc manner) definitely can be used in a single process / single node env as well, it is not only useful for distributed engines.
We still several pipeline breaking operators, e.g. sorts and hash join build side, which basically needs to load the full input before it can make progress.
Currently this could be (and already is) done ad-hoc inside each operator, but a (small) framework to mark pipeline breaking stages and dynamically re-optimizing the plan sounds like a more principled way of doing things.
Some obvious examples include:
- Join reordering
- Push down dynamic filters (currently it is done inside operators in a non-generalizable way)
- Update statistics (for join, aggregations or parallel merge)
Reacted by Andrew Lamb and Brent GardnerThe way I see it is, we could add the rules that make sense for running on a single machine, distributed engines could extend this with extra rules (just like the current logical / physical optimizer rule list).
Reacted by Oleks V and Brent GardnerIf we could somehow migrate Ballista to use datafusion-distributed I think that would also be amazing (and would remove all doubt about moving datafusion-distributed into the core datafusion crate)
I would be eager to learn what would be the benefit of doing so?
Reacted by Brent Gardner and Andy GroveRather than prescribing AQE implementation should the focus be on providing building blocks to implement AQE planners?
From perspective of ballista AQE the hard part was to make "event driven" core which would re-plan stages based on statistics collected at the stage boundaries.
Significant discrepancy from datafusion execution is that datafusion makes quite a lot decisions before plan starts (eg. partitioning, or join implementation), AQE would benefit from deferring such decisions for later in the flow. https://www.cs.cmu.edu/~15721-f24/papers/AQP_in_Lakehouse.pdf mentions that spark converts parts of logical to physical plan when needed rather than on execution start, I believe similar approach would simplify AQE implementation (ballista implements DynamicJoin to defer some partitioning and join decisions, which i see as a hack)
IMHO having a planner which could act on execution events would be the first step
Reacted by Daniël Heres, Andy Grove and Adrian Garcia BadaraccoReacted by Daniël Heres, Brent Gardner, Andy Grove and Andrew LambIf we could somehow migrate Ballista to use datafusion-distributed I think that would also be amazing (and would remove all doubt about moving datafusion-distributed into the core datafusion crate)
I would be eager to learn what would be the benefit of doing so?
in my (biased) view it would
- Validate that datafusion-distributed is general enough that it can be used with an existing system like Ballista
- It would potentially let us pool development efforts
To be clear I am not saying it is required or even a good idea -- I was mostly dreaming that if it was to come to pass it would be cool
The way I see it is, we could add the rules that make sense for running on a single machine, distributed engines could extend this with extra rules (just like the current logical / physical optimizer rule list).
Spark does the same thing
Speaking to AQE overall it would be nice to extend EXPLAIN opportunities to show what was the physical plan before and after AQE, or it should be something else, as AQE should display the plan modified during the execution.
Also if speaking to Spark the AQE rebuilds plan and stages bottom up when
AdaptiveSparkPlanExec, typically after shuffle IFF child nodes materialized(DF currently doesn't track it). The rest of the plan rebuilt, reevalulated and optimizedextend EXPLAIN opportunities to show what was the physical plan before and after AQE
I hit this problem in the PR SLT. The plan has been rewritten, but I have no way to show it as it is currently being executed. I'm open to suggestions.
IFF child nodes materialized(DF currently doesn't track it)
The PR calls
prime()on stages, and knows if what point they have hadPhysicalPlanNode::execute()called on them. We could make this more explicit with a.has_started()method or something.edit: I'd just like to highlight, because things tend to get lost in the noise, this isn't just an abstract issue I filed. There is a real, working, generalizable example of AQE in DF PR attached to this issue. Please, have your LLM take a look :) (to the general audience, not @comphead in particular)
Reacted by Oleks VIf we could somehow migrate Ballista to use datafusion-distributed I think that would also be amazing (and would remove all doubt about moving datafusion-distributed into the core datafusion crate)
I would be eager to learn what would be the benefit of doing so?
in my (biased) view it would
- Validate that datafusion-distributed is general enough that it can be used with an existing system like Ballista
- It would potentially let us pool development efforts
To be clear I am not saying it is required or even a good idea -- I was mostly dreaming that if it was to come to pass it would be cool
I wouldn't really support directly adding datafusion-distributed as a "vetted" distributed solution / library.
datafusion-distributed has a significantly different architecture / execution model (streaming network exchange vs shuffle, task vs complete stage, etc.) than others so effectively using it as a library isn't really a feasible option.
Rather, I would support building / extending the building blocks necessary for distributed engines like datafusion-distributed, ballista, comet, etc. that can be used by those.
If that means some of the code can be re-used/donated or inspired from either datafusion-distributed/ballista/etc... I am all for that, but I think it should be done piece-by-piece.
in my (biased) view it would
- Validate that datafusion-distributed is general enough that it can be used with an existing system like Ballista
- It would potentially let us pool development efforts
To be clear I am not saying it is required or even a good idea -- I was mostly dreaming that if it was to come to pass it would be cool
One thing to note here is that the goal of
datafusion-distributedis to be general for end users, allowing them to build their own distributed engine for their own datasources, for example:datafusionextension points allow people to provide their own data sources, and customize how they are executeddatafusion-distributedextension points allow people to customize how those data sources should be distributed, what workers should be involved, how work is split, where should network boundaries be placed, etc... all while maintaining an almost exact 1:1 execution model as normaldatafusion, just happening that some nodes move data over the wire rather than in-memory.
So, the design goal of
datafusion-distributedis to satisfy end users directly, which I'm not sure if Ballista would fit in that definition. There do are success stories of other systems, like ParadeDB, that are currently built on top ofdatafusion-distributedthough, but I don't think it's in the same position as Ballista.The market gap
datafusion-distributedfills is also different from Ballista:- Customizable In-memory engines: Trino is the closest here, but it's not nearly as customizable as the
datafusion+datafusion-distributedcombo, and we are confident that a system built on topdatafusion-distributedis faster and cheaper to run than their Trino equivalent. - Customizable disk-based shuffling engines: Ballista would fall in this bucket, but there are already really big players doing a very good job in this space (Spark, Comet), so
datafusion-distributedis not interested in competing here.
Datafusion's first concern, rightly, was not AQE; thus, trying to add one at this moment may be more of a hack rather than a systematic approach to it.
I agree with @Dandandan comment #23194 (comment) making a plan-observation/runtime-adaptation framework could be a step in right direction, and would help to make some separation of concerns between operators and planner(s). IMHO, making a case how core datafusion would benefit from this would be a good first step.
I'd suggest we focus on requirements/gaps rather than prescribing solutions. Also, it would be great if we refrain from speculative comparisons and claims, as they do not really help.
Reacted by Gabriel- Post-shuffle partition coalescing — collapse many tiny partitions (over-estimated repartition) into a smaller number of right-sized ones. Needs per-partition row counts.
- Adaptive aggregation strategy — pick between hash and sort aggregation, or between
PartialandSingle, based on observed cardinality.
To add to these points, one "adaptive execution" method we had seen the need for was this: #20847
Being able to dynamically decide whether we need to partitioned aggregate/partitioned join is something that DataFusion doesn't naturally support. The output partitions (and their nature, like unknown/hash) are fixed at plan time.
Reacted by Andrew LambI have a small PR (#25798) proposing the stage-boundary interface as a first, standalone piece: an execution-plan capability to prime input partitions, observe when each has completed, and release output. It includes a test/example implementation and examples for deterministic pause/resume and caller-controlled stage admission. It does not add AQE orchestration, runtime optimizer rules, or statistics. I’m sharing it here because it may be a useful core seam for the broader directions in this issue. Feedback on whether this interface belongs in core—and whether its per-partition readiness/release contract is the right scope—would be appreciated.
Reacted by Brent Gardner
Is your feature request related to a problem or challenge?
DataFusion makes all execution-plan choices at planning time from static statistics, which are frequently
InexactorAbsent(e.g. after aGROUP BY, a selective filter, or anything involving derived columns). A wide set of optimizations is blocked by this — the planner has to guess at runtime cardinality and a bad guess costs orders of magnitude.Concrete decisions that need runtime stats:
PartialandSingle, based on observed cardinality.Each of these needs the same primitive: read whatever runtime stat is relevant to the decision from a completed sub-plan, then re-run a planning decision against the up-to-date plan.
Describe the solution you'd like
PoC in #23167.
The shape of the solution:
StageBoundaryBuffer— anExecutionPlanoperator that materializes its input. It is a pipeline breaker by construction, so it's a safe synchronization point where runtime stats become observable. No need to find a natural breaker inside the operator being adapted.runtime_*stat methods onExecutionPlan, each populated by the buffer's drain task as data flows through. The PoC shipsruntime_row_count(used for the build-side swap). The same pattern extends toruntime_partition_extrema(in-progress for parallel sort / range repartitioning), per-key distribution sketches (for skew handling), bloom-filter snapshots (for dynamic filter pushdown), and so on. Long-term these consolidate intopartition_statisticsreturningPrecision::Exactonce the breaker has completed.InsertHashJoinBoundaries(the first, in the PoC) wraps each HashJoin input. Future adaptive optimizations each get their own targeted insertion rule (InsertWindowBoundaries,InsertRangePartitioningBoundaries, etc.). The rules don't need to know about each other.RuntimeOptimizerExec— single always-on wrapper at the plan root. On each stage-completion event, fires registered runtime rules against the current plan, replaces the plan with each rule's output, then releases the stage's boundaries so downstream consumers can drain. Critically, the runtime rule can rewrite any part of the plan above the boundary — different partition counts, different repartition strategies (hash → range), different join layouts — so dynamic partitioning falls out of the same machinery; the buffer just stabilizes the input.RuntimeRuletrait with identical signature toPhysicalOptimizerRule::optimize. They live in separate traits today only becausephysical-plancannot depend onphysical-optimizer(the dependency runs the other way). End-state: a single trait. Existing rules become runtime-aware by reading boundary state from the plan tree they receive —JoinSelectionitself can run against a post-breaker plan with no logic change, only fresher numbers.RepartitionExec's spawn model).MemoryPoolpressure, spill to disk if available; otherwise release the boundary early and let the rule see "overflowed" — the query still runs, we just lose the adaptive decision (no worse than the pre-AQE baseline).The PoC demonstrates the build-side-swap case end-to-end. Every other item in the problem list is some combination of (a) one ~30-line targeted insertion rule, (b) one runtime rule that reads the relevant stat method, and (c) whichever
runtime_*stat method that decision needs.Describe alternatives you've considered
Inexactestimates. Bad estimates produce slow queries. No path to win on the optimizations listed above.partition_statisticsas data flows through; rules poll and re-decide. Possible but requires contract changes across every operator and a polling-based fire mechanism. The boundary approach gives the same epistemic guarantee with a single new operator and event-driven firing.Additional context
PoC PR: #23167. The runtime swap on a real query is observable via
RUST_LOG=info:Spark AQE for context: https://spark.apache.org/docs/latest/sql-performance-tuning.html#adaptive-query-execution