Conversation
28bbcab to
0e2e532
Compare
adf49d1 to
8e06826
Compare
8e06826 to
d18a86c
Compare
|
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 |
|
Indeed. I actually wonder if we should even reuse One foundational invariant of |
|
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 |
|
@stuhood this is comething I am going to draft toda na dadd to the epic as well 🙇 |
Awesome, works for me.
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? |
d18a86c to
e5900d5
Compare
e5900d5 to
0a62dad
Compare
Stacked atop #734.
Summary
Adds support for
Partitioning::Rangeacross stage boundaries (NetworkShuffleExec), scales range split points to consumer task counts during stage preparation, and fixes partition count tracking and multi-consumer scaling inNetworkCoalesceExec.Previously:
inject_network_boundariesonly checked forPartitioning::Hash. When DataFusion produced a range repartition, noNetworkShuffleExecwas injected, causing rows to remain within producer tasks and dropping data that needed to cross task boundaries.NetworkCoalesceExecassumed 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)). Whenconsumer_tasks > 1, this inflated advertised partition counts, causing downstream operators to query non-existent partition streams.RangePartitioningeven 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:
NetworkShuffleExecforPartitioning::Range.RangePartitioning::scale.consumer_taskstoNetworkCoalesceExecand scale properties byinput_tasks.div_ceil(consumer_tasks).max(1)during stage preparation (try_from_stage,with_input_stage, andwith_new_children).consumer_tasks(tag 5) toNetworkCoalesceExecProtoinsrc/codec/distributed_codec.rs, defaulting to 1 for backwards compatibility.src/execution_plans/common.rsandsrc/stage.rs.UnknownPartitioninginsrc/events/defaults/file_scan_config.rswhen a task receives fewer file groups than the range partition count.Testing and snapshots:
RangePartitionedTableWrapper::scanintests/range_partitioning.rsto group files contiguously by path rather than using round-robin grouping, preserving sorted range boundaries.tests/range_partitioning.rsandtests/dynamic_filtering/partitioned_join.rsto reflect the execution fixes.NetworkCoalesceExecscaling, codec roundtrip withconsumer_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 tostuhood.branch-55-range-scalingonparadedb/datafusionto incorporate apache/datafusion#24766 before it lands upstream.