Skip to content

Returning Worker Pool Results in Submission Order

Concurrent workers finish items out of order, but many consumers need results in the order items were submitted: rows written to a file, events applied to a state machine, a paginated export. Sorting at the end works only when you can wait for everything. Streaming results in order while the pool keeps running needs a reorder buffer: hold results that arrive early until every earlier result has been emitted. Its size is the hidden cost. Tested with 1,000 items on 8 workers, where 1% of items took 500 ms and the rest 10 ms: an unbounded buffer emitted results in order with the same 2.19 s total as an unordered run, but held up to 286 results at once behind the stragglers; bounding the window to 64 items in flight capped the buffer at 64 and stretched the run to 5.42 s, because each straggler stalled the window. This guide builds both and shows how to pick the bound.

Prerequisites

1. Tag items with a sequence number

Order is only recoverable if each item carries its position. Number items as they are submitted and carry the number through to the result:

import asyncio


async def run_pool(items, handle, workers: int = 8):
    work: asyncio.Queue = asyncio.Queue()
    results: asyncio.Queue = asyncio.Queue()

    async def worker():
        while True:
            seq, item = await work.get()
            try:
                results.put_nowait((seq, await handle(item)))
            except Exception as exc:
                results.put_nowait((seq, exc))           # keep the slot even on failure
            finally:
                work.task_done()

    for seq, item in enumerate(items):
        work.put_nowait((seq, item))
    tasks = [asyncio.create_task(worker()) for _ in range(workers)]
    return results, tasks

Failures must occupy their slot. If a failed item produced no result, the reorder buffer would wait forever for its sequence number and every later result would be stuck behind it. Emitting the exception as the result lets the consumer decide what a failed position means.

Verify: inject one failing item; the stream still completes, with the exception at that item's position.

2. Emit through a reorder buffer

The consumer keeps a dict keyed by sequence number and emits the next expected one whenever it is present:

async def in_order(results: asyncio.Queue, total: int):
    buffer: dict[int, object] = {}
    next_seq = 0
    while next_seq < total:
        seq, value = await results.get()
        buffer[seq] = value
        while next_seq in buffer:
            yield buffer.pop(next_seq)
            next_seq += 1

Measured: output identical to the input order, total time 2.19 s — the same as an unordered run, because workers never wait on the consumer. The first result was emitted after 10 ms. The cost appeared in memory: when item 37 took 500 ms, the 280-odd items behind it finished and sat in the buffer, peaking at 286 entries. With results of a few kilobytes that is negligible; with result sets of megabytes, it is a memory spike proportional to the slowest item's duration times throughput.

Verify: track len(buffer) and confirm its peak is acceptable for your result size.

A slow item holds later results in the buffer A sequence of 6 messages between 3 participants. A slow item holds later results in the buffer workers reorder buffer consumer result 0, result 1 emit 0, 1 results 3, 4, 5 (2 still running) hold 3, 4, 5 result 2 (slow) emit 2, 3, 4, 5 Everything that finishes behind a straggler waits in memory until the straggler is done.

3. Bound the window to bound memory

To cap memory, cap how far ahead of the consumer the pool may run. A semaphore released as results are emitted (not completed) limits items in flight plus items buffered:

async def ordered_bounded(items, handle, workers: int = 8, window: int = 64):
    slots = asyncio.Semaphore(window)
    work: asyncio.Queue = asyncio.Queue()
    results: asyncio.Queue = asyncio.Queue()

    async def feeder():
        for seq, item in enumerate(items):
            await slots.acquire()                       # wait until the consumer catches up
            await work.put((seq, item))

    ...                                                 # workers as in step 1
    buffer, next_seq = {}, 0
    while next_seq < len(items):
        seq, value = await results.get()
        buffer[seq] = value
        while next_seq in buffer:
            yield buffer.pop(next_seq)
            next_seq += 1
            slots.release()                             # one emitted, one more may start

Measured with a window of 64: buffer peak exactly 64, but total time rose from 2.19 s to 5.42 s. Each 500 ms straggler stalled the window: once 63 later items were done and buffered, no new item could start until the straggler finished, so workers sat idle. This is head-of-line blocking, and it is the fundamental trade-off of ordered output with bounded memory.

Verify: run with your real latency distribution at a few window sizes and plot total time against peak buffer.

Ordered output: buffer size versus total time A grid of 3 rows by 4 columns. Ordered output: buffer size versus total time mode total time peak buffered first result unordered 2.20 s 0 as soon as any finishes ordered, unbounded buffer 2.19 s 286 10 ms ordered, window 64 5.42 s 64 10 ms Memory and throughput trade against each other only when there are stragglers.

4. Size the window from the latency tail

The window that avoids stalls is about throughput × the slowest item's duration: everything that can finish while one straggler runs. For 8 workers on 10 ms items (about 800 per second) and 500 ms stragglers, that is roughly 400 — consistent with the unbounded run's peak of 286. Choose:

def window_for(throughput_per_s: float, p999_item_s: float, max_buffer: int) -> int:
    no_stall = int(throughput_per_s * p999_item_s) + 1
    return min(no_stall, max_buffer)

If no_stall fits in memory, use it and get full throughput. If not, the window is a memory cap and you accept stalls — or attack the tail directly: a per-item timeout turns a 500 ms straggler into a fast failure that occupies its slot (see setting per-item timeouts in worker pools), and hedging can cut it further.

Verify: at the chosen window, workers are rarely idle — measure idle time per worker — while the buffer stays within budget.

5. Order per key instead of globally when you can

Global order is often stronger than needed. If results only need to be ordered per key — per account, per document — items for different keys can be emitted independently, and a straggler for one key does not hold back others:

from collections import defaultdict


async def in_order_per_key(results: asyncio.Queue, counts: dict[str, int]):
    buffers: dict[str, dict[int, object]] = defaultdict(dict)
    next_seq: dict[str, int] = defaultdict(int)
    remaining = sum(counts.values())
    while remaining:
        (key, seq), value = await results.get()
        buffers[key][seq] = value
        while next_seq[key] in buffers[key]:
            yield key, buffers[key].pop(next_seq[key])
            next_seq[key] += 1
            remaining -= 1

Each key has its own sequence and buffer, so head-of-line blocking is confined to the key that owns the straggler. When per-key order is all that is required, routing each key to one worker removes the buffer entirely, as in processing queue items in order per key.

Verify: a straggler for key A does not delay emission of key B's results.

How much ordering does the consumer need? A decision on What order must results come out in with 3 outcomes. How much ordering does the consumer need? What order must results come out in? any order emit as completed no buffer per key buffer per key or route by key stragglers stay local global input order reorder buffer window from the tail Ask for the weakest ordering the consumer can live with; it is the cheapest to provide.

Verification

Ordered output is correct when:

  • Output order equals input order in tests with randomised latencies.
  • Failures occupy their slot, so a failed item never blocks the stream.
  • Buffer size is bounded by a window chosen from measured throughput and tail latency.
  • Workers are rarely idle at the chosen window, or the stall cost is understood and accepted.

Diagnostic Hook: export the reorder buffer's size and the age of the item the consumer is waiting for. A buffer pinned at the window size with a waiting-item age far above the median item duration is a straggler stalling the pipeline; logging that item's key identifies whether the slowness is data-dependent.

Pitfalls & edge cases

  • Dropping failed items. The buffer waits forever for their sequence numbers.
  • Releasing the window on completion instead of emission. The buffer becomes unbounded again.
  • Global ordering where per-key would do. One slow key stalls everything.
  • Large results in an unbounded buffer. Memory spikes in proportion to the tail latency.

Frequently Asked Questions

How do I get results from an asyncio worker pool in input order?

Tag each item with a sequence number, have workers return the number with the result, and emit results through a reorder buffer that holds early arrivals until every earlier position has been emitted.

How large does a reorder buffer get?

About throughput multiplied by the duration of the slowest in-flight item. With 8 workers on 10 ms items and 500 ms stragglers, it peaked at 286 results in testing.

How do I bound the memory of ordered output?

Limit how far the pool may run ahead of the consumer with a semaphore released when results are emitted. This caps the buffer but causes head-of-line stalls: a window of 64 turned a 2.2 s run into 5.4 s.

What happens if an item fails in an ordered pool?

It must still produce a result for its position — typically the exception — or the reorder buffer will wait forever for that sequence number.