Skip to content

Choosing chunksize for ProcessPoolExecutor Work

Every item sent to a process pool pays a fixed toll: pickle the function reference and arguments, push them through a pipe, unpickle in the worker, and do the same for the result on the way back. When the work per item is small, the toll dominates. Measured on Python 3.14 with eight worker processes, mapping a trivial function over 100,000 items took 7.70 s with chunksize=1, 0.73 s with 10, 0.18 s with 100 and 0.07 s with 1,000. From asyncio the situation is worse, because loop.run_in_executor submits one item at a time and has no chunksize at all: 20,000 items of about 50 µs of work each took 1.93 s submitted individually and 0.17 s in batches of 100 — slower individually than simply running them serially on one core (0.99 s). This guide picks a chunk size and shows how to batch from async code.

Prerequisites

1. Measure the per-item overhead first

The right chunk size depends on the ratio between per-item work and per-task overhead, so measure both on your workload:

import time
from concurrent.futures import ProcessPoolExecutor


def small(x: int) -> int:
    return x * x


def run(fn, n: int, chunksize: int) -> float:
    with ProcessPoolExecutor(max_workers=8) as pool:
        list(pool.map(fn, range(10)))                     # start the workers first
        t = time.perf_counter()
        list(pool.map(fn, range(n), chunksize=chunksize))
        return time.perf_counter() - t


if __name__ == "__main__":
    for cs in (1, 10, 100, 1000):
        print(cs, round(run(small, 100_000, cs), 3))

For the trivial function, chunksize=1 spent about 77 µs per item on overhead — far more than the multiplication itself. A function doing roughly 50 µs of work showed the same shape less dramatically: 1.36 s at chunksize 1, 0.23 s at 10, and 0.18 s from 100 upwards. Once chunks are large enough that per-chunk overhead is small next to per-chunk work, larger chunks stop helping.

Verify: your own curve flattens at some chunk size; that knee, or a little above it, is your starting point.

100,000 trivial items across 8 processes 4 horizontal bars comparing chunksize 1 with the others. 100,000 trivial items across 8 processes chunksize 1 7.70 s chunksize 10 0.73 s chunksize 100 0.18 s chunksize 1,000 0.07 s Python 3.14, 8 worker processes, x*x per item; pool started before timing. With tiny items the pipe and pickle cost is the whole runtime; chunking amortises it away.

2. Pick a chunk size from the work, not the item count

A useful rule: make each chunk take at least a few milliseconds of worker time, and make sure there are at least several chunks per worker so the load balances.

import math


def choose_chunksize(n_items: int, workers: int, est_item_s: float,
                     min_chunk_s: float = 0.005, chunks_per_worker: int = 4) -> int:
    by_time = math.ceil(min_chunk_s / max(est_item_s, 1e-9))         # amortise overhead
    by_balance = max(1, n_items // (workers * chunks_per_worker))    # keep all workers busy
    return max(1, min(by_time, by_balance))

For 20,000 items of 50 µs on 8 workers, by_time gives 100 and by_balance gives 625, so 100 — which matches the knee measured above. Too few chunks is its own problem: with 8 chunks for 8 workers, one slow chunk leaves seven workers idle at the end. The default chunksize=1 in ProcessPoolExecutor.map is safe for heavy items and pathological for light ones; never leave it at the default for fine-grained work.

Verify: compare the chosen value's runtime with neighbours at half and double; it should be within a few percent of the best.

3. Batch from asyncio yourself

loop.run_in_executor(pool, fn, item) is the asyncio entry point to a process pool, and it submits exactly one call. Batch explicitly: send a list per call and process it in the worker:

import asyncio
from concurrent.futures import ProcessPoolExecutor


def score(x: int) -> int:
    return sum(i * x for i in range(2000))


def score_batch(xs: list[int]) -> list[int]:          # runs in the worker process
    return [score(x) for x in xs]


async def score_all(pool: ProcessPoolExecutor, items: list[int], size: int = 100) -> list[int]:
    loop = asyncio.get_running_loop()
    chunks = [items[i:i + size] for i in range(0, len(items), size)]
    results = await asyncio.gather(*(loop.run_in_executor(pool, score_batch, c) for c in chunks))
    return [r for chunk in results for r in chunk]


if __name__ == "__main__":
    async def main():
        with ProcessPoolExecutor(max_workers=8) as pool:
            print(len(await score_all(pool, list(range(20_000)))))
    asyncio.run(main())

Measured on 20,000 items: 1.93 s individually, 0.17 s with batches of 100, 0.17 s with 500, 0.19 s with 2,500 (only eight batches, so worse balance). The batch function must be a module-level function so it can be pickled by reference.

Verify: the number of run_in_executor calls equals the number of chunks, not the number of items.

20,000 items of ~50 µs from asyncio 4 horizontal bars comparing one run_in_executor per item with the others. 20,000 items of ~50 µs from asyncio one run_in_executor per item 1.93 s serial, no pool at all 0.99 s batches of 2,500 (8 calls) 0.19 s batches of 100 (200 calls) 0.17 s 8 worker processes, Python 3.14; batches processed by a list comprehension in the worker. Unbatched, the pool was slower than doing the work on one core.

4. Stream results instead of waiting for all chunks

Gathering every chunk before using any result holds all results in memory and delays the first. With asyncio.as_completed, consume chunks as they finish:

async def score_streaming(pool, items, size: int = 100):
    loop = asyncio.get_running_loop()
    futures = {
        loop.run_in_executor(pool, score_batch, items[i:i + size]): i
        for i in range(0, len(items), size)
    }
    for fut in asyncio.as_completed(futures):
        start, results = await fut
        yield start, results

That requires the batch function to return its starting offset alongside results, so the consumer can place them. For very large inputs, also bound how many chunks are in flight with a semaphore, so the parent process does not pickle the entire input up front — the bounded fan-out pattern from processing results in completion order with as_completed.

Verify: the first results are available after roughly one chunk's processing time, and parent-process memory does not scale with input size.

5. Watch out for the start method and chunk payloads

Two operational details change the numbers:

  • Start method. On Python 3.14, Linux uses forkserver by default: workers start from a clean server process, not a copy of the parent, so module-level state in the parent is not inherited and every worker re-imports your modules. Startup is slower and the __main__ guard is required — the first version of the benchmark above crashed with "An attempt has been made to start a new process before the current process has finished its bootstrapping phase" until the guard was added.
  • Chunk payload size. A chunk of 1,000 large objects is pickled into one message; if items are big, the payload, not the call count, becomes the cost. Measure bytes per chunk, and move bulk data through shared memory instead, as described in sharing large arrays with shared memory.
import multiprocessing as mp

if __name__ == "__main__":
    ctx = mp.get_context("forkserver")              # explicit, same on every version
    with ProcessPoolExecutor(max_workers=8, mp_context=ctx) as pool:
        ...

Setting the context explicitly makes behaviour identical across 3.13 (fork default) and 3.14 (forkserver default).

Verify: the pool behaves the same on every Python version you support, and chunk payloads stay in the kilobytes-to-low-megabytes range.

How should this work reach the process pool? A decision on How much work is one item with 3 outcomes. How should this work reach the process pool? How much work is one item? tens of ms or more one call per item overhead is noise micro to a few ms batch to ~5 ms per call amortise the toll big data per item shared memory do not pickle it The pool's per-call toll is roughly fixed; size calls so it is a small fraction of each.

Verification

Chunking is right when:

  • Per-call overhead is a small fraction of per-call work, measured.
  • Each worker gets several chunks, so load balances at the end.
  • asyncio code batches explicitly rather than calling run_in_executor per item.
  • The start method is explicit, and every entry point has a __main__ guard.

Diagnostic Hook: record per-batch wall time and items per batch as a histogram in the parent. Batches far shorter than 5 ms mean overhead dominates and the size should grow; a long tail of slow batches at the end of a run means too few chunks for good balancing. Compare total pool time against a serial baseline occasionally — a pool slower than serial is the clearest sign of over-fine submission.

Pitfalls & edge cases

  • Default chunksize=1 with tiny items. The pool can be slower than doing the work on one core.
  • Lambdas or closures as the batch function. They cannot be pickled; use module-level functions.
  • Huge chunks. Few chunks per worker leaves workers idle at the end and delays first results.
  • Missing __main__ guard on 3.14. Forkserver re-imports the main module and fails at pool start.

Frequently Asked Questions

What chunksize should I use with ProcessPoolExecutor.map?

Large enough that each chunk takes at least a few milliseconds of work, small enough that each worker gets several chunks. In testing, 100,000 trivial items took 7.7 s at chunksize 1 and 0.07 s at 1,000.

Does run_in_executor support chunksize?

No. It submits one call per invocation. Split the input into lists yourself and submit a module-level function that processes a whole list, which took 20,000 small items from 1.93 s to 0.17 s.

Why is my process pool slower than running serially?

Per-item overhead — pickling arguments and results and passing them through pipes — exceeds the work per item. Batch the items so each call does more work.

Why does Python 3.14 need if name == 'main' for process pools on Linux?

The default start method on Linux changed from fork to forkserver in 3.14. Worker processes import the main module, so pool creation must not run at import time.