Repository navigation
Supporting analytics over large amounts of pre‑partitioned data #24438
Description
Activity
take
@alamb and @gene-bordegaray : This is the feature request we chatted earlier
Reacted by Gene Bordegaray@xavlee : feel free to take this ticket and use the branch or create a new branch
take
untake
- added a commit that references this issue
on Aug 20, 2026 Copy some comments from #24497
I see now that the idea here is that somehow the external system knows the data is not sorted but is non overlapping
I think this is going to be really hard to manage / ensure through the plan -- specifically I think the challenge will be to ensure that every operator properly reports if the property is preserved or not
I am not sure this is something we want to complicate datafusion with at this time as it
- Is complicated (see above)
- I am not sure how many other systems have data organized in this way / could take advantage of it
We are already working on Range Partitioning (which is another form of distribution aware planning) which is also complicated, but I think has a clearly demonstrated wide demand.
Let's see if we can find others who need this before we try and code it
Thanks, @alamb for outlining the complication concerns and asking for concrete use cases.
Use case:
We’re entering an AI‑driven era with extremely high‑volume telemetry and analysis workloads. On the same dataset, we now face two extremes:- Interactive queries — require low latency. We partition data finely so queries run in parallel. Both DataFusion and Distributed DataFusion handle this well, and we scale perfectly.
- Analytic queries — must finish within acceptable latency on very large datasets. Here, the number of partitions far exceeds the number of cores (even across many workers), so we must merge partitions. We’re fine with sequential execution within a merged partition as long as it completes in reasonable time. But merging drops the sort order, which causes two major issues:
i .Loss of streaming aggregation, hurting downstream efficiency and latency.
ii. For high‑cardinality workloads, hash tables become huge, increasing resource pressure and risking OOMs under concurrency.
We now see the core idea: the external system can indicate that data is unsorted but non‑overlapping. This is similar to other metadata properties (e.g., sort order) that come from the external system.
On complexity:
Yes, it’s complicated — and we’re happy to design and support this (thanks @gene-bordegaray and @xavlee and future Datadog contributors) with input from the DataFusion community. We see this as an opportunity to push DataFusion to the next level for AI workloads, with Datadog acting as a pioneer in providing real use cases.Let us know your thoughts, and whether you’d like us to bring in other users of range partitioning to contribute.
We now see the core idea: the external system can indicate that data is unsorted but non‑overlapping. This is similar to other metadata properties (e.g., sort order) that come from the external system.
Yes, I think is likely a common pattern (and the use case you describe of "more partitions than cores" certainly happens for other users -- like we have it at Influx for example
Maybe that is the usecase we can highlight -- "supporting analytics over large amounts of pre-partitioned data" 🤔 It sounds pretty neat and not something that I am aware has been written about before
Reacted by Nga Tran4 remaining items
Thanks @suremarc for the quick analysis, tests, and suggestions — really appreciate the effort and the momentum it gives this feature request.
If you GROUP BY time_partition, key, date_bin(time) then the aggregate becomes fully streaming.
We considered this approach and would use it if necessary, but it’s still sub‑optimal compared with the proposal in this ticket:
- It requires adding a property that data is sorted on
partition_key, key, timestampinstead of key, timestamp. - Users naturally write
GROUP BY key, date_bin(time), and we would need to rewrite it to includepartition_key. - Grouping on three keys is more expensive than grouping on two.
- Any operators before the group‑by may break the required sort order, leaving us with the cost of a three‑key group‑by without the benefit.
The proposed feature feels much more natural:
- It avoids all the drawbacks above.
- It fits cleanly into the optimizer: treating non‑overlapping data on group‑by keys as a first‑class property can help not only this case but future optimizations as well.
@alamb — this might even be a small research‑driven insight applied to a real‑world workload. It’s not trivial, but meaningful optimizations rarely are.
@xavlee, @gene-bordegaray, @jayshrivastava and I have aligned on a reasonable design. Xavier already has four draft PRs: one for tests demonstrating the case, two for plumbing, and one integrating everything. Gene, Jayant, and I will review carefully so Xavier can address all comments before we bring it to the committers.
We hope this feature helps push DataFusion to the next level — showcasing a strong query engine ready for critical workloads across the world.
- It requires adding a property that data is sorted on
Thanks @suremarc for looping me into the issue.
Yes, I think this generalization would be useful for Massive. We have tables partitioned by day or month and hash-binned by key, with each logical time partition sorted by
(key, time). When Atlas combines many logical partitions into one DataFusion execution partition, the effective ordering becomes(time_partition, key, time). A query such asGROUP BY key, date_bin(time)therefore loses a globally usable(key, time)ordering, even though its groups remain non-overlapping when the time buckets do not cross partition boundaries. This is a direct use case for the original proposal.We also have many materialization queries that read time-ordered data and compute OHLC rollups grouped by key and time bucket using ordered
FIRST_VALUEandLAST_VALUE. Those inputs are not necessarily fully group-contiguous because keys may interleave, but knowing that time is ordered when the key is fixed would allow the aggregate to retain only the current bucket per key.Today we compensate with additional execution partitions or hash-repartitioned aggregation, both of which have costs. From Massive's perspective, the conditional/per-logical-partition ordering model looks broadly useful for materialization, window functions, and potentially streaming execution.
I'll follow up the issue and the related PRs
Reacted by Nga TranAwesome — thanks, @xudong963. I appreciate the quick follow‑up and look forward to your review and input.
@xudong963, @gene-bordegaray and @jayshrivastava: Here are the four PRs to review, in order.
Reacted by Gene Bordegaray and xudong.wThank you for the in depth explanation and discussion. I caight up on this and review the first testing PR. Thank you @xavlee
A trend I am seeing across the project is the need for more partitioning properties to properly be represented and taken avantage of. This largely started with preserving pre partitioned data to eliminate shuffles, then new range partitioning use cases, and more like this are popping up in the repo. Given us at Datadog are one of the main drivers of these efforts, it seems that others using Datafusion at scale also see the use case for this.
I think it would be ineresting to start a discussion about where the community wants to see partitioning head and how DF should expose properties that ir can take advantage of. Hopefully others here agree with that and I can start this discussion and link relevant work going on around this. 🙇
Reacted by xavlee, Dan Harris and Nga TranAll we need upstream is a property that can be propagated to AggregateExec to indicate this.
Shouldn't this property already be inferrable from column statistics? If the
[min_key, max_key]intervals are disjoint across input partitions, and we haveOrderingEquivalanceClasswith[key ASC/DESC, date_bin(time) ASC/DESC]then we should be able to stream the aggregation without additional machinery right?All we need upstream is a property that can be propagated to AggregateExec to indicate this.
Shouldn't this property already be inferrable from column statistics? If the
[min_key, max_key]intervals are disjoint across input partitions, and we haveOrderingEquivalanceClasswith[key ASC/DESC, date_bin(time) ASC/DESC]then we should be able to stream the aggregation without additional machinery right?This will work. And datafusion checks for this right now and will stream. But this property streatches further then that use case. The main idea is concatting sorted runs can lose the global sort but still preserve group contiguousness (wow first time I have used that word and surprisinly it is a real word 😆 ):
run 1: AAA BBB run 2: CCC DDD concat (2 -> 1): CCC DDD AAA BBBEvery group is contiguous. Thusm it can be emitted incrementally but the combined stream cannot advertise sorted order.
Exact, disjoint per-run stats could infer this property, but merged/inexact stats cannot. So stats should be one way to derive the property but we still need a property to represent and propagate the guarantees to agg.Reacted by Dan Harris and Nga TranThanks @gene-bordegaray. A discussion sounds great
I was speaking with people at VLDB this week, and we were discussing how to represent this notion with our existing infrastructure.
Goetz Graefe suggested a formalized version of what I think @xudong963 and @gene-bordegaray describe as 'group-contiguous' in #24438 (comment) and #24438 (comment).
The idea is to model this as an extension of the existing
SortProperties--ASC,DESCand (the new)GROUPED. This would then naturally fit into all the existing sort analysis infrastructure. Here is how these relate, via exampleASC: rows are sorted by ascendingmonthmonth | value ------+------ Jan | 100 Jan | 42 Feb | 77 Feb | 12 Mar | 55 Mar | 91DESC: rows are sorted by descendingmonthmonth | value ------+------ Mar | 55 Mar | 91 Feb | 77 Feb | 12 Jan | 100 Jan | 42GROUPED(new): all rows with the samemonthare contiguous, but the months themselves appear in no particular ordermonth | value ------+------ Feb | 77 Feb | 12 Mar | 55 Mar | 91 Jan | 100 Jan | 42GROUPEDis strictly weaker thanASC/DESC(any sorted input is also grouped), but I think it is sufficient for the various streaming operations we have in DataFusion (e.g. window functions and grouping) . For example, for streaming aggregation: once the value ofmonthchanges, the aggregator knows that group is complete and can emit it.One challenge with this approach would be that SortProperties I think is in arrow-rs and updating it / wrapping it in DataFusion would be fairly invasive / a breaking API cange
The changes from @xavlee in #24698 is less disruptive, but is also basically a special case
Reacted by Nga Tran and Samuel ArnoldThanks @alamb — the GROUPED model is a nice way to formalize what we’ve been calling group-contiguous.
Two quick clarifications, then a question:
SortPropertiesalready lives in DataFusion (datafusion-expr-common), not arrow-rs. The arrow type isSortOptions(ASC/DESC+ nulls).GROUPEDis not a total order, so we would not put it onPhysicalSortExpr/SortOptions(that would touch EnforceSorting, SMJ, proto, etc.).If we go the
GROUPEDroute, the implementation I’d propose is:- Keep real sorts as they are (
PhysicalSortExpr+oeq_class). - Add a sibling
geq_classonEquivalenceProperties: a lex tuple of exprs that are contiguous but not ordered. Every existing ordering implies a grouping (drop ASC/DESC). - Add
SortProperties::Groupedonly for expression propagation (date_binalready copies its input’s sort property). - Sources (esp.
FileScanConfig, when concat dropsoutput_orderingbut ranges are disjoint) advertise the weaker grouping. AggregateExecuses that for group completion. We would still want Xavier’s refactor: use group clustering terminology in aggregate execution #24697 split (GroupCompletionMode≠InputOrderMode). feat: add Grouped equivalence properties #24698’s sidecar onPlanPropertieswould not be needed.
#24698 is narrower and less work;
GROUPEDonEquivalencePropertiesis the more general property.Which would you like us to pursue — implement
GROUPEDas above, or continue with Xavier’s #24698 / #24497 stack?- Keep real sorts as they are (
@alamb and I just synced in person and agreed to move forward with the new
GROUPEDapproach above.@xavlee: would you be able to implement this on top of #24697?
FYI: @xudong963 and @gene-bordegaray
Thanks for the discussion. I'll update the design starting from #24697 with the
GROUPEDapproach.The stack has been reviewed and approved by @gene-bordegaray.
- test: cover unsorted contiguous groups in one partition #24737
- refactor: use group clustering terminology in aggregate execution #24697
- feat: add Grouped equivalence properties #24698
- feat: stream aggregates over grouped input #24497
cc @xudong963 and @alamb for review.
Reacted by Nga Tran, Gene Bordegaray and xudong.w- added a commit that references this issue
on Sep 28, 2026 - added a commit that references this issue
on Oct 11, 2026
Overview
After discussing with @alamb in the comments, I updated the subject to supporting analytics over large amounts of pre‑partitioned data. One concrete ask for that use case (from the original subject) is:
@xavlee: If this becomes a larger epic, you may want to add well‑defined subtasks
Problem statement
Our data is range‑partitioned on two dimensions,
timeandkey, and each file is sorted by(key, time).This layout allows us to execute fully streaming the query below very efficiently, as shown in the plan below. Each of our partitions can be considered as one file-group and mapped directly to a DataFusion partition.
The challenge arises when a query needs to scan many more data partitions than the number of CPU cores, which is also the default for
target_partitions. Since we all know it’s not recommended to settarget_partitionsfar above the CPU count, we’re forced to merge many of our data partitions into a single DataFusion partition. Once we do that, we lose the(key, time)sort order, which meansAggregateExeccan no longer stream the data.Describe the solution you'd like
Looking at the query plan below with the partitioning described above, we can see that even though each DataFusion partition (stream) is not sorted, the execution is still fully streaming. This works because the data across partitions does not overlap on the grouping keys (key, date_bin(..)).
If we introduce a new property that tells AggregateExec the input is non‑overlapping on the group‑by keys, then it can safely execute in a fully streaming fashion even without a global sort order.
We’ll handle the merging of many data partitions into a single DataFusion partition on our side and ensure the merged data remains non‑overlapping. All we need upstream is a property that can be propagated to AggregateExec to indicate this.
Update:
Describe alternatives you've considered
No response
Additional context
No response