Skip to content

Supporting analytics over large amounts of pre‑partitioned data #24438

Description

@NGA-TRAN

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:

  • Support streaming aggregates when partitions are unsorted but non‑overlapping streaming aggregates

@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, time and key, and each file is sorted by (key, time).

Image

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.

SELECT key, date_bin(...), sum(...)
FROM    my_table
GROUP BY  key, date_bin(...)
Image

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 set target_partitions far 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 means AggregateExec can 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.

Image

Update:

  1. For the first solution, we may only need to propagate this property from the datasource through ProjectionExec (if necessary) and into the first AggregateExec. Whether it should propagate beyond that point requires more thought, but that’s outside the scope of this feature request. We can simply stop propagation there.
  2. If it passes through ProjectionExec, we may initially support only monotonic functions such as date_bin, since that won’t affect correctness or behavior.

Describe alternatives you've considered

No response

Additional context

No response

Activity

  1. added theissue type on Aug 17, 2026
  2. NGA-TRAN commented on Aug 17, 2026

    @NGA-TRAN
    ContributorAuthor

    take

  3. NGA-TRAN commented on Aug 17, 2026

    @NGA-TRAN
    ContributorAuthor

    @alamb and @gene-bordegaray : This is the feature request we chatted earlier

  4. NGA-TRAN commented on Aug 17, 2026

    @NGA-TRAN
    ContributorAuthor

    @xavlee : feel free to take this ticket and use the branch or create a new branch

  5. xavlee commented on Aug 17, 2026

    @xavlee
    Contributor

    take

  6. NGA-TRAN commented on Aug 17, 2026

    @NGA-TRAN
    ContributorAuthor

    untake

  7. alamb commented on Aug 25, 2026

    @alamb
    Contributor

    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

    1. Is complicated (see above)
    2. 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

  8. gene-bordegaray commented on Aug 25, 2026

    @gene-bordegaray
    Contributor

    I can / want to help support this effort (with #22497 as sent to the dev list). I know this is not another need since @NGA-TRAN and I are coming from the same place but would be more than willing to support reviews and dev work for this 🙇

  9. NGA-TRAN commented on Aug 25, 2026

    @NGA-TRAN
    ContributorAuthor

    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:

    1. 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.
    2. 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.

  10. alamb commented on Aug 26, 2026

    @alamb
    Contributor

    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

  11. 4 remaining items

  12. NGA-TRAN commented on Aug 30, 2026

    @NGA-TRAN
    ContributorAuthor

    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:

    1. It requires adding a property that data is sorted on partition_key, key, timestamp instead of key, timestamp.
    2. Users naturally write GROUP BY key, date_bin(time), and we would need to rewrite it to include partition_key.
    3. Grouping on three keys is more expensive than grouping on two.
    4. 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:

    1. It avoids all the drawbacks above.
    2. 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.

  13. xudong963 commented on Aug 31, 2026

    @xudong963
    Member

    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 as GROUP 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_VALUE and LAST_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

  14. NGA-TRAN commented on Aug 31, 2026

    @NGA-TRAN
    ContributorAuthor

    Awesome — thanks, @xudong963. I appreciate the quick follow‑up and look forward to your review and input.

  15. NGA-TRAN commented on Aug 31, 2026

    @NGA-TRAN
    ContributorAuthor
  16. gene-bordegaray commented on Aug 31, 2026

    @gene-bordegaray
    Contributor

    Thank 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. 🙇

  17. thinkharderdev commented on Aug 31, 2026

    @thinkharderdev
    Contributor

    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 have OrderingEquivalanceClass with [key ASC/DESC, date_bin(time) ASC/DESC] then we should be able to stream the aggregation without additional machinery right?

  18. gene-bordegaray commented on Aug 31, 2026

    @gene-bordegaray
    Contributor

    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 have OrderingEquivalanceClass with [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 BBB
    

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

  19. NGA-TRAN commented on Sep 3, 2026

    @NGA-TRAN
    ContributorAuthor

    Thanks @gene-bordegaray. A discussion sounds great

  20. alamb commented on Sep 3, 2026

    @alamb
    Contributor

    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, DESC and (the new) GROUPED. This would then naturally fit into all the existing sort analysis infrastructure. Here is how these relate, via example

    ASC: rows are sorted by ascending month

    month | value
    ------+------
    Jan   |   100
    Jan   |    42
    Feb   |    77
    Feb   |    12
    Mar   |    55
    Mar   |    91
    

    DESC: rows are sorted by descending month

    month | value
    ------+------
    Mar   |    55
    Mar   |    91
    Feb   |    77
    Feb   |    12
    Jan   |   100
    Jan   |    42
    

    GROUPED (new): all rows with the same month are contiguous, but the months themselves appear in no particular order

    month | value
    ------+------
    Feb   |    77
    Feb   |    12
    Mar   |    55
    Mar   |    91
    Jan   |   100
    Jan   |    42
    

    GROUPED is strictly weaker than ASC / 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 of month changes, 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

  21. NGA-TRAN commented on Sep 3, 2026

    @NGA-TRAN
    ContributorAuthor

    Thanks @alamb — the GROUPED model is a nice way to formalize what we’ve been calling group-contiguous.

    Two quick clarifications, then a question:

    SortProperties already lives in DataFusion (datafusion-expr-common), not arrow-rs. The arrow type is SortOptions (ASC/DESC + nulls). GROUPED is not a total order, so we would not put it on PhysicalSortExpr / SortOptions (that would touch EnforceSorting, SMJ, proto, etc.).

    If we go the GROUPED route, the implementation I’d propose is:

    1. Keep real sorts as they are (PhysicalSortExpr + oeq_class).
    2. Add a sibling geq_class on EquivalenceProperties: a lex tuple of exprs that are contiguous but not ordered. Every existing ordering implies a grouping (drop ASC/DESC).
    3. Add SortProperties::Grouped only for expression propagation (date_bin already copies its input’s sort property).
    4. Sources (esp. FileScanConfig, when concat drops output_ordering but ranges are disjoint) advertise the weaker grouping.
    5. AggregateExec uses 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 on PlanProperties would not be needed.

    #24698 is narrower and less work; GROUPED on EquivalenceProperties is the more general property.

    Which would you like us to pursue — implement GROUPED as above, or continue with Xavier’s #24698 / #24497 stack?

  22. NGA-TRAN commented on Sep 3, 2026

    @NGA-TRAN
    ContributorAuthor

    @alamb and I just synced in person and agreed to move forward with the new GROUPED approach above.

    @xavlee: would you be able to implement this on top of #24697?

    FYI: @xudong963 and @gene-bordegaray

  23. xavlee commented on Sep 3, 2026

    @xavlee
    Contributor

    Thanks for the discussion. I'll update the design starting from #24697 with the GROUPED approach.

  24. xavlee commented on Sep 23, 2026

    @xavlee
    Contributor
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

enhancementNew feature or request

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions