Skip to content

feat: Support Range partitioning across stage boundaries - #730

Open
stuhood wants to merge 3 commits into
datafusion-contrib:mainfrom
paradedb:stuhood.three-way-range-upstream
Open

stuhood wants to merge 3 commits into
datafusion-contrib:mainfrom
paradedb:stuhood.three-way-range-upstream

Conversation

@stuhood

@stuhood stuhood commented Sep 16, 2026 •

Copy link
Copy Markdown
Contributor

Stacked atop #734.

Summary

Adds support for Partitioning::Range across stage boundaries (NetworkShuffleExec), scales range split points to consumer task counts during stage preparation, and fixes partition count tracking and multi-consumer scaling in NetworkCoalesceExec.

Previously:

  • inject_network_boundaries only checked for Partitioning::Hash. When DataFusion produced a range repartition, no NetworkShuffleExec was injected, causing rows to remain within producer tasks and dropping data that needed to cross task boundaries.
  • NetworkCoalesceExec assumed 1 consumer task and scaled output partitions by total input tasks (local.tasks) rather than per-consumer group size (input_tasks.div_ceil(consumer_tasks)). When consumer_tasks > 1, this inflated advertised partition counts, causing downstream operators to query non-existent partition streams.
  • Range-partitioned leaf scans scaled across tasks retained their global RangePartitioning even when a task was assigned only a subset of file groups, causing partition count mismatches downstream.

Additionally, closes #628 by completing the execution support for the benchmarks and tests added in #734.

Changes

  • Planning and execution:

    • Inject NetworkShuffleExec for Partitioning::Range.
    • Scale range split points across consumer tasks using RangePartitioning::scale.
    • Add consumer_tasks to NetworkCoalesceExec and scale properties by input_tasks.div_ceil(consumer_tasks).max(1) during stage preparation (try_from_stage, with_input_stage, and with_new_children).
    • Add consumer_tasks (tag 5) to NetworkCoalesceExecProto in src/codec/distributed_codec.rs, defaulting to 1 for backwards compatibility.
    • Split shuffle fan-out scaling from coalesce partition scaling in src/execution_plans/common.rs and src/stage.rs.
    • Downgrade leaf scan partitioning to UnknownPartitioning in src/events/defaults/file_scan_config.rs when a task receives fewer file groups than the range partition count.
  • Testing and snapshots:

    • Update RangePartitionedTableWrapper::scan in tests/range_partitioning.rs to group files contiguously by path rather than using round-robin grouping, preserving sorted range boundaries.
    • Update test plan and result snapshots in tests/range_partitioning.rs and tests/dynamic_filtering/partitioned_join.rs to reflect the execution fixes.
    • Add unit tests for NetworkCoalesceExec scaling, codec roundtrip with consumer_tasks, and leaf scan range downgrade.

Review Notes

Stacked atop #734: this PR contains the execution engine changes and updates the goldens introduced in #734 to pass.

Note: Workspace dependencies in Cargo.toml (datafusion, datafusion-proto, datafusion-cli) temporarily point to stuhood.branch-55-range-scaling on paradedb/datafusion to incorporate apache/datafusion#24766 before it lands upstream.

@stuhood
stuhood force-pushed the stuhood.three-way-range-upstream branch from 28bbcab to 0e2e532 Compare September 17, 2026 03:04
@stuhood stuhood changed the title Support Range partitioning across stage boundaries feat: Support Range partitioning across stage boundaries Sep 17, 2026
@stuhood
stuhood force-pushed the stuhood.three-way-range-upstream branch 2 times, most recently from adf49d1 to 8e06826 Compare September 17, 2026 23:23
@stuhood
stuhood marked this pull request as ready for review September 18, 2026 04:40
@stuhood
stuhood force-pushed the stuhood.three-way-range-upstream branch from 8e06826 to d18a86c Compare September 18, 2026 04:41
@gene-bordegaray

Copy link
Copy Markdown
Collaborator

I haven't started reviewing but @Rich-T-kid this will definitely corss cu with the shuffle work and we should coordinate how this should land / be integrated.

Also cc: @gabotechs

@gabotechs

Copy link
Copy Markdown
Collaborator

Indeed. I actually wonder if we should even reuse NetworkShuffleExec for the Range repartitioning.

One foundational invariant of NetworkShuffleExec is that the connection between producer and consumer stages is all-to-all, meaning that each individual consumer needs to pull data from all the producers below. However, I think this is not necessarily true for a range re-partition boundary, a consumer only needs to pull data from whatever producers contain SplitPoints overlaping the ones the consumer expects.

@gene-bordegaray

Copy link
Copy Markdown
Collaborator

Ya this is what @jayshrivastava had brought up as well: on a range we can greatly reduce the amount of connections between producers and consumers on range. I would be very interested if this coul dbe smootly integrated upstream in the repartition and DFD could be near to oblivious of the arcitecture change.

I would think this would be the route to go since single-node would also be very happy about each input partition having to fanout to less output partitions

@gene-bordegaray

Copy link
Copy Markdown
Collaborator

@stuhood this is comething I am going to draft toda na dadd to the epic as well 🙇

@stuhood

stuhood commented Sep 18, 2026 •

Copy link
Copy Markdown
Contributor Author

@stuhood this is comething I am going to draft toda na dadd to the epic as well 🙇

Awesome, works for me.


Indeed. I actually wonder if we should even reuse NetworkShuffleExec for the Range repartitioning.
..
I would think this would be the route to go since single-node would also be very happy about each input partition having to fanout to less output partitions

I'm definitely in favor of fixing this in the right place.

But downstream, we also need an interim solution. If we are able to get this cleaned up to a "working and landable but not optimal" state, then that would be my preference, with TODOs pointed at upstream tickets.

Does that seem feasible?

stuhood added a commit to paradedb/datafusion-distributed that referenced this pull request Sep 18, 2026
stuhood added a commit to paradedb/datafusion-distributed that referenced this pull request Sep 18, 2026
@stuhood
stuhood force-pushed the stuhood.three-way-range-upstream branch from d18a86c to e5900d5 Compare September 18, 2026 18:44
@stuhood
stuhood force-pushed the stuhood.three-way-range-upstream branch from e5900d5 to 0a62dad Compare September 18, 2026 19:05
stuhood added a commit to paradedb/datafusion-distributed that referenced this pull request Sep 18, 2026
stuhood added a commit to paradedb/datafusion-distributed that referenced this pull request Sep 19, 2026
stuhood added a commit to paradedb/datafusion-distributed that referenced this pull request Sep 19, 2026
philippemnoel pushed a commit to paradedb/datafusion-distributed that referenced this pull request Sep 20, 2026
stuhood pushed a commit to paradedb/datafusion-distributed that referenced this pull request Sep 22, 2026
barbarj pushed a commit to paradedb/datafusion-distributed that referenced this pull request Sep 24, 2026
philippemnoel pushed a commit to paradedb/datafusion-distributed that referenced this pull request Sep 25, 2026

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

tests and benchmarks on range partitioned data

3 participants