Skip to content

Improve internal worker parallelism support #23174

Description

@2010YOUY01

Is your feature request related to a problem or challenge?

Original discussion from @alamb : #23026 (comment)

Summary: DataFusion mostly uses repartition-based parallelism today, but at some point we need to introduce intra-partition parallelism, and we have to do that carefully for performance.

The existing SortExec has already exposed some internal worker parallelism: it is possible to have a large number of concurrent workers for local sorting.

  • let streams = std::mem::take(&mut self.in_mem_batches)
    .into_iter()
    .map(|batch| {
    let metrics = self.metrics.baseline.intermediate();
    let reservation = self
    .reservation
    .split(get_reserved_bytes_for_record_batch(&batch)?);
    let input = self.sort_batch_stream(batch, &metrics, reservation)?;
    Ok(spawn_buffered(input, 1))

This issue explains the background and proposes some ideas for improving this support.

What is internal worker parallelism

DataFusion mostly uses repartition based parallelism, each partition has independent data, and we use 1 CPU core to process one partition, here is a parallel aggregation query example:

Image

For certain workloads, the assumptions for repartition are not ideal, here are 3 motivating examples

Motivating Example 1: memory pressure case

Let's say we're doing a large sort (data size >> memory), on a machine with 32 cores, 64GB memory. The default setting will execute it with 32 global partitions, and with each partition a classic external sort algorithm is executed (local sort, spill disk, and finally read back and sort-preserving merge)
The issue is that per-partition memory budget is low, the spilling might create smaller sorted runs, and the end-to-end execution requires extra spills, reading back, and merging smaller files.
A more ideal plan shape is:

  • the scanner keep 32 partitions, since they're not memory intensive and scalable with number of CPU cores
  • SortExec shrink to 8 partitions, with 4 internal workers per partition. This gives each partition more memory budget to proceed easier; also there are efficient algorithms to parallelize sort and sort-preserving merge within partition.
Image

Motivating Example 2: segment-tree based parallelism in window functions

The window query in the figure is impossible to parallelize with repartition, because it assumes data independence among partitions, and the query has one global partition, and window frame changes every row.
At the meantime, there is a very parallel algorithm if we can allow shared memory among partitions:

Then the ideal query shape become

(any downstream exec)
-- RepartitionExec(round-robin on batch, input_partitions=1, otuput_partitions=32)
---- WindowExec(partition=1, internal_parallelism=32)
------ CoalescePartitionExec(input_partitions=32, output_partitions=1)
-------- CsvExec(partition=32, internal_parallelism=1)
Image

Motivating Example 3

This PR from @Dandandan seems also tries to introduce intra partition parallelism

Describe the solution you'd like

  • Establish conventions for intra-partition parallelism. For example, each execution plan stage may want to maintain a similar total concurrency level: partition_count * internal_workers_per_partition. See the CsvExec + WindowExec example above.
  • Improve Explain output for internal parallelism. The existing single-partition SortExec can still use internal parallelism, but the plan currently looks serial, which makes it hard to inspect potential performance issues.

Describe alternatives you've considered

No response

Additional context

No response

Activity

  1. alamb commented on Jun 25, 2026

    @alamb
    Contributor

    Summary: DataFusion mostly uses repartition-based parallelism today, but at some point we need to introduce intra-partition parallelism, and we have to do that carefully for performance.

    Another way to describe this, that I prefer is "partition based parallelism" -- basically DataFusion tries to create a plan that will use target_partition number of partitions, which will then use that many cores during execution.

    I think there are some counter examples that don't live up to this exactly (e.g. wide fan in Unions and SortPreservingMerge) but otherwise it is largely the same and works well. Part of the reason DataFusion is so fast for the classic Scan-filter-aggregate is this partition parallelism (see my talk at TokioConf about Using Tokio for CPU-Bound Tasks (Works Really Well) TokioConf 2026 ( slides, and recording)

  2. alamb commented on Jun 25, 2026

    @alamb
    Contributor

    I think the property that the partitioning / parallelism is explicit in the plan is a core one in DataFusion

    Another way we could consider modeling the usecases above (e.g. a single partition in a WindowFunction) is to keep the multiple partitions, but internally use shared state, similar to how RepartitionExec or JoinExec, and FileInputStream. That way we keep the parallelism tied to the plan's structure (partition_count) but the streams executing each partition can dynamically adapt during plan time to better use resources

    Concretely, maybe instead of

    (any downstream exec)
    -- RepartitionExec(round-robin on batch, input_partitions=1, otuput_partitions=32)
    ---- WindowExec(partition=1, internal_parallelism=32)
    ------ CoalescePartitionExec(input_partitions=32, output_partitions=1)
    -------- CsvExec(partition=32, internal_parallelism=1)
    

    We had something like

    (any downstream exec)
    ---- WindowExec(partition=32)
    -------- CsvExec(partition=32)
    

    And then internally within WindowExec(partition=32, internal_parallelism=32) it knew enough to coalesce the data into a single partition, and then split the work across the multiple partition streams (with a segment tree or whatever)

    I don't think this is fundamentally different than a single exec with internal paralleism, but I think it keeps task/cpi model the same

  3. 2010YOUY01 commented on Jun 26, 2026

    @2010YOUY01
    ContributorAuthor

    I think the property that the partitioning / parallelism is explicit in the plan is a core one in DataFusion

    Another way we could consider modeling the usecases above (e.g. a single partition in a WindowFunction) is to keep the multiple partitions, but internally use shared state, similar to how RepartitionExec or JoinExec, and FileInputStream. That way we keep the parallelism tied to the plan's structure (partition_count) but the streams executing each partition can dynamically adapt during plan time to better use resources

    Concretely, maybe instead of

    (any downstream exec)
    -- RepartitionExec(round-robin on batch, input_partitions=1, otuput_partitions=32)
    ---- WindowExec(partition=1, internal_parallelism=32)
    ------ CoalescePartitionExec(input_partitions=32, output_partitions=1)
    -------- CsvExec(partition=32, internal_parallelism=1)
    

    We had something like

    (any downstream exec)
    ---- WindowExec(partition=32)
    -------- CsvExec(partition=32)
    

    And then internally within WindowExec(partition=32, internal_parallelism=32) it knew enough to coalesce the data into a single partition, and then split the work across the multiple partition streams (with a segment tree or whatever)

    I don't think this is fundamentally different than a single exec with internal paralleism, but I think it keeps task/cpi model the same

    I think this simplified model is a good idea in single-node cases. My concern is that, for distributed use cases, shared-state parallelism usually isn't easily applicable, so we may want to explicitly separate partition parallelism from internal parallelism to make the adaptation easier.

    cc @gabotechs who has been working on datafusion-distributed

  4. alamb commented on Jun 27, 2026

    @alamb
    Contributor

    My concern is that, for distributed use cases, shared-state parallelism usually isn't easily applicable,

    This is an excellent point -- if we want to use the same paralleization model for both local (cross core) and distributed (cross node) it would be much harder (e.g. internally the ExecutionPlan nodes would have to communicate acros nodes somehow or something) which sounds pretty complicated.

    On the other hand, maybe it is ok if we had two different approaches for cross core parallelism and cross node parallelism as the constraints are somewhat different

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions