Skip to content

Streaming Results from a Process Pool

Offloading CPU work from asyncio to a ProcessPoolExecutor is usually written as "submit everything, gather everything". That returns nothing until the slowest task finishes, and it keeps every task's input alive in the parent until then. When there are many tasks, results should stream out as they complete, and inputs should be submitted only as fast as workers can take them. Measured on Python 3.14 with 8 worker processes and 2,000 tasks, each taking 0.3–12 ms of CPU and a distinct 100 KB input: gather over all futures delivered its first result after 1,057 ms — when everything was done — with 197 MiB of inputs held in the parent. asyncio.as_completed over the same 2,000 futures delivered the first result after 92 ms but still held 197 MiB. Keeping a window of 16 tasks in flight and submitting a new one as each finished delivered the first result after 2.2 ms, finished all 2,000 in 0.95 s, and peaked at 1.8 MiB. multiprocessing.Pool.imap_unordered, bridged to the loop, took 21 ms, 0.85 s and 1.0 MiB. This guide builds the bounded stream.

Prerequisites

1. See what gather costs on many tasks

The common pattern submits every task at once and waits for all of them:

async def process_all(pool, items):
    loop = asyncio.get_running_loop()
    futures = [loop.run_in_executor(pool, crunch, item) for item in items]
    return await asyncio.gather(*futures)

Measured with 2,000 tasks of 100 KB each: the first result was available after 1,057 ms, which was also when the last one was, and the parent's traced memory peaked at 197 MiB — every input stayed referenced by its pending work item until the whole batch completed. For a report that needs every result at once, waiting is unavoidable, but the memory is not; for anything that can act on results individually — writing rows, sending progress, forwarding to another stage — gathering throws away a second of latency per batch and makes memory proportional to batch size.

Verify: you know whether your caller needs all results together or can process them one by one.

2. Use as_completed for early results, knowing its memory cost

asyncio.as_completed yields futures in completion order, so results arrive as soon as each task finishes:

async def stream_as_completed(pool, items):
    loop = asyncio.get_running_loop()
    futures = [loop.run_in_executor(pool, crunch, item) for item in items]
    for next_done in asyncio.as_completed(futures):
        yield await next_done

Measured: the first result arrived after 92 ms instead of 1,057, and all 2,000 were done in 1.07 s. Memory did not improve — 197 MiB — because every task, and its input, was still submitted up front. For a few dozen tasks this is the simplest correct answer. For thousands, or unbounded inputs such as rows read from a file, it is the same memory problem as gather with better latency. Results arrive out of order; tag each input with an index if the caller needs to match them, as shown in processing results in completion order with as_completed.

Verify: the number of tasks submitted at once is small enough that holding all their inputs is acceptable.

2,000 CPU tasks with distinct 100 KB inputs, 8 workers A grid of 4 rows by 4 columns. 2,000 CPU tasks with distinct 100 KB inputs, 8 workers pattern first result all 2,000 peak memory in parent gather over all futures 1,057 ms 1.06 s 197.4 MiB as_completed over all futures 92 ms 1.07 s 197.1 MiB bounded window, 16 in flight 2.2 ms 0.95 s 1.8 MiB Pool.imap_unordered, chunksize 4 21 ms 0.85 s 1.0 MiB Python 3.14, forkserver workers; memory traced with tracemalloc.

3. Keep a bounded window of tasks in flight

To bound memory as well as latency, submit only a few tasks more than there are workers, and submit the next one each time a result comes back:

async def stream_bounded(pool, items, window: int):
    loop = asyncio.get_running_loop()
    pending: set[asyncio.Future] = set()
    source = iter(items)

    def submit_next() -> None:
        item = next(source, None)
        if item is not None:
            pending.add(loop.run_in_executor(pool, crunch, item))

    for _ in range(window):
        submit_next()
    while pending:
        done, _ = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED)
        for future in done:
            pending.discard(future)
            submit_next()
            yield future.result()

Measured with a window of 16 — twice the worker count, so each worker always has its next task queued: the first result after 2.2 ms, all 2,000 in 0.95 s, and 1.8 MiB of peak memory in the parent. Inputs are created as the window advances, so items can be a generator over a file or a database cursor of any size. The window is the knob: at least the number of workers to keep them all busy, a little more to hide the round trip, and far less than the batch size.

Verify: peak memory in the parent stays flat as the number of tasks grows, and all workers stay busy.

The bounded window as results arrive A sequence of 6 messages between 4 participants. The bounded window as results arrive consumer stream_bounded process pool input source take 16 inputs submit 16 tasks first completed (2.2 ms) yield result take 1 input submit 1 task: 16 in flight Inputs are read only as fast as results come back.

4. Or use multiprocessing.Pool.imap_unordered

multiprocessing.Pool.imap_unordered implements the same idea inside the pool: it consumes an input iterator lazily and yields results as they complete, batching small tasks with chunksize. It is synchronous, so run the iteration in a thread and pass results to the loop:

async def stream_imap(items, workers: int = 8, chunksize: int = 4):
    loop = asyncio.get_running_loop()
    queue: asyncio.Queue = asyncio.Queue()
    ctx = mp.get_context("forkserver")

    def run() -> None:
        with ctx.Pool(workers) as pool:
            for result in pool.imap_unordered(crunch, items, chunksize=chunksize):
                loop.call_soon_threadsafe(queue.put_nowait, result)
        loop.call_soon_threadsafe(queue.put_nowait, None)

    runner = loop.run_in_executor(None, run)
    while (result := await queue.get()) is not None:
        yield result
    await runner

Measured with chunksize=4: first result after 21 ms, all 2,000 in 0.85 s — the fastest total, because chunking cut the per-task overhead — and 1.0 MiB of peak memory. imap_unordered reads ahead from the input iterator to keep its task queue full, so for truly unbounded inputs, watch its memory; and a slow consumer lets results pile up in the asyncio queue, which a maxsize would turn into backpressure on the thread. Chunking is discussed in more depth in choosing chunksize for process pool map.

Verify: results stream as they complete, and the bridge thread finishes when the input is exhausted.

5. Handle failures and cancellation in the stream

A stream of results also needs a policy for failed tasks and for a consumer that stops early. In the bounded window, an exception surfaces from future.result() for that task only, and the rest can continue; if the consumer stops iterating, the tasks still in flight should be cancelled:

async def stream_bounded_safe(pool, items, window: int):
    loop = asyncio.get_running_loop()
    pending: set[asyncio.Future] = set()
    source = iter(items)
    try:
        for _ in range(window):
            if (item := next(source, None)) is not None:
                pending.add(loop.run_in_executor(pool, crunch, item))
        while pending:
            done, _ = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED)
            for future in done:
                pending.discard(future)
                if (item := next(source, None)) is not None:
                    pending.add(loop.run_in_executor(pool, crunch, item))
                if future.exception() is not None:
                    yield ("error", future.exception())
                else:
                    yield ("ok", future.result())
    finally:
        for future in pending:
            future.cancel()                      # tasks not yet started are dropped

Cancelling an executor future removes the task if a worker has not picked it up; a task already running continues to completion in its process, as described in cancelling long-running work in a process pool. Consume the generator with contextlib.aclosing so the finally runs as soon as the consumer stops.

Verify: a consumer that breaks out early leaves at most window tasks running, and a failing task does not end the stream.

How should results come back from the pool? A decision on How many tasks, and how are results used with 4 outcomes. How should results come back from the pool? How many tasks, and how are results used? few, all needed together gather first result = last few, used one by one as_completed 92 ms to first many or unbounded bounded window, ~2x workers 2.2 ms, 1.8 MiB many tiny tasks imap_unordered + chunksize 0.85 s total Submitting everything up front costs memory proportional to the batch.

Verification

Process-pool results are streamed well when:

  • The first result is available long before the last, for callers that can use it.
  • Tasks in flight are bounded, so parent memory does not grow with batch size.
  • Workers stay busy, with a window of at least the worker count.
  • Failures are reported per task, and pending tasks are cancelled when the consumer stops.

Diagnostic Hook: track the executor's pending work items — futures submitted but not completed — alongside parent memory. A count that jumps to the batch size at the start of each run and drains slowly is the submit-everything pattern; under a bounded window, it should hover at the window size for the whole run.

Pitfalls & edge cases

  • gather over thousands of tasks. Measured: first result after 1,057 ms, 197 MiB held.
  • as_completed over everything. Measured: early results, the same 197 MiB.
  • A window smaller than the worker count. Workers sit idle.
  • Abandoning the stream without cancelling. Queued tasks keep running for nobody.

Frequently Asked Questions

How do I get results from a ProcessPoolExecutor as they finish in asyncio?

Keep a bounded set of run_in_executor futures, wait with asyncio.wait(..., FIRST_COMPLETED), yield each result and submit the next input. With 16 in flight, the first of 2,000 results came after 2.2 ms and the parent used 1.8 MiB.

Why does gather over a process pool use so much memory?

All tasks are submitted at once, and each input stays referenced until the batch completes: 2,000 inputs of 100 KB held 197 MiB in testing.

Does asyncio.as_completed solve the memory problem?

No: it gives early results (92 ms here) but still submits everything up front, so memory matched gather.

Should I use multiprocessing.Pool.imap_unordered instead?

It streams lazily with chunking and was fastest overall (0.85 s) with 1.0 MiB in the parent, but it is synchronous and needs a thread to bridge into asyncio.