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¶
- Python 3.11+.
- Process pools from asyncio, from offloading CPU work with loop.run_in_executor.
- The topic overview, CPU-Bound Task Offloading.
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.
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.
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.
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¶
gatherover thousands of tasks. Measured: first result after 1,057 ms, 197 MiB held.as_completedover 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.
Related¶
- CPU-Bound Task Offloading — up to the topic overview.
- Recycling process pool workers — keeping long-running pools healthy.
- Concurrent Execution & Worker Patterns — the section overview.