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¶
- Python 3.11+, stdlib only.
- Worker pools, from building an async worker pool with TaskGroup.
- The iterator equivalent, from writing async itertools helpers.
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.
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.
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.
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.
Related¶
- Worker Pool Implementations — up to the topic overview.
- Preserving order across concurrent pipeline stages — the same problem across several stages.
- Concurrent Execution & Worker Patterns — the section overview.