Preserving Order Across Concurrent Pipeline Stages¶
Concurrent stages are what make a pipeline fast, and they scramble order: with eight workers per stage, the item that entered first is rarely the first to leave. When the consumer needs input order — rows appended to a file, events applied to a state machine, a log replayed — the order has to be restored. Doing it after every stage is the instinctive approach and is rarely needed. Measured on a two-stage pipeline of 1,000 items, eight workers per stage, with about 1% of items taking 200 ms in the first stage: re-sequencing after each stage held at most 190 items in the first buffer and 10 in the second; re-sequencing once at the sink held at most 192 — the same memory and the same 1.23–1.25 s, with one buffer instead of two. This guide tags items once, re-sequences once, and covers the cases where an intermediate re-sequence is genuinely required.
Prerequisites¶
- Python 3.11+, stdlib only.
- Reorder buffers, from returning worker pool results in submission order.
- Staged pipelines, from building a staged async pipeline with bounded queues.
1. Tag items with a sequence number at the source¶
Order can only be restored if every item carries its original position. Assign it once, where items enter the pipeline, and carry it through every stage unchanged:
import asyncio
from dataclasses import dataclass
from typing import Any
@dataclass(slots=True)
class Envelope:
seq: int
value: Any
async def source(items, out: asyncio.Queue) -> int:
n = 0
async for item in items:
await out.put(Envelope(n, item))
n += 1
return n # total, needed by the re-sequencer
Stages operate on envelope.value and pass the envelope on. Keeping the sequence number in a wrapper rather than mixing it into the payload means stage code cannot accidentally drop or change it. If a stage filters items out, it must still emit something for that sequence number — a tombstone envelope — or the re-sequencer will wait forever for the gap.
Verify: every envelope that leaves the last stage has a sequence number, and the set of numbers is exactly range(total).
2. Re-sequence once, at the sink¶
A single reorder buffer before the order-sensitive consumer restores input order regardless of how many concurrent stages came before:
async def resequence(inbox: asyncio.Queue, total: int):
buffer: dict[int, Any] = {}
next_seq = 0
while next_seq < total:
env = await inbox.get()
buffer[env.seq] = env.value
inbox.task_done()
while next_seq in buffer:
yield buffer.pop(next_seq)
next_seq += 1
Measured on two concurrent stages: one buffer at the sink peaked at 192 items, while buffers after each stage peaked at 190 and 10 — the same total memory. The run time was the same too (1.23 s against 1.25 s). An intermediate buffer adds a stage boundary, a task and a place for bugs, and buys nothing unless the next stage needs ordered input. The buffer's size is set by how far the fastest items get ahead of the slowest — mostly by the slow first-stage items, wherever the buffer sits.
Verify: output order equals input order, and the buffer's peak is close to throughput × the slowest in-flight item's latency.
3. Re-sequence mid-pipeline only for stateful stages¶
An intermediate re-sequence is required when a stage's processing depends on order: computing running totals, deduplicating against the previous record, applying a diff to the previous state, detecting session boundaries. Such a stage must see items in order, so it gets a re-sequencer in front and runs with a single worker — or with one worker per key, if order only matters per key:
async def stateful_running_total(inbox: asyncio.Queue, outbox: asyncio.Queue, total: int):
running = 0
async for value in resequence(inbox, total): # ordered input
running += value.amount
await outbox.put(Envelope(value.seq, {**value.__dict__, "running": running}))
Everything before it can be concurrent; the stateful stage itself is sequential, because its correctness depends on processing one item after another. Keep it cheap, and push heavy work into the concurrent stages before it. If only per-key order matters — running totals per account — route each key to its own sequential worker instead, as in processing queue items in order per key.
Verify: the stateful stage's output matches a sequential reference implementation run on the same input.
4. Bound the buffer with a window, and know the cost¶
A reorder buffer's size is unbounded in principle: one item stuck for minutes lets everything behind it pile up. Bound it by limiting how far the source may run ahead of the re-sequencer — a semaphore acquired at tagging and released when an item is emitted in order:
async def bounded_source(items, out: asyncio.Queue, window: asyncio.Semaphore) -> int:
n = 0
async for item in items:
await window.acquire() # released by the re-sequencer on emit
await out.put(Envelope(n, item))
n += 1
return n
The trade-off is the one measured for worker pools: a window caps memory, but a slow item then stalls the whole pipeline once the window fills behind it, idling workers. Size the window from throughput times the latency of your slow tail, and attack the tail itself with per-item timeouts, as in setting per-item timeouts in worker pools. A timed-out item still needs a tombstone at its sequence number so ordering continues.
Verify: with the window in place, peak buffer equals the window, and throughput drops only when stragglers exceed the window's time budget.
5. Handle failures and filters with tombstones¶
A sequence number that never arrives blocks the re-sequencer forever. Every stage must emit exactly one envelope per input envelope — the transformed value, a filtered marker, or an error:
SKIP = object()
def stage_worker(fn):
async def worker(inbox: asyncio.Queue, outbox: asyncio.Queue):
while True:
env = await inbox.get()
try:
result = await fn(env.value)
await outbox.put(Envelope(env.seq, SKIP if result is None else result))
except Exception as exc:
await outbox.put(Envelope(env.seq, exc)) # keep the slot
finally:
inbox.task_done()
return worker
The sink decides what to do with SKIP (ignore it) and with exceptions (log, dead-letter, or stop). Because every sequence number arrives exactly once, ordering cannot stall on a lost item, and the output's position-for-position correspondence with the input is preserved — useful when the output must line up with the input, such as writing results back next to the rows they came from.
Verify: inject filters and failures at random; the re-sequencer always completes, and the sink sees exactly total envelopes.
Verification¶
Ordering is preserved correctly when:
- Items are tagged once, at the source, and the tag survives every stage.
- Re-sequencing happens only in front of order-sensitive consumers.
- Every stage emits exactly one envelope per input, using tombstones for filtered and failed items.
- Buffer memory is bounded by a window sized from the latency tail.
Diagnostic Hook: export the re-sequencer's buffer size and the sequence number it is waiting for, alongside that item's age. A buffer growing while the awaited sequence number stays fixed is a straggler; its age tells you whether a per-item timeout would have caught it. If the awaited number never arrives at all, a stage dropped an envelope — the bug the tombstone rule exists to prevent.
Pitfalls & edge cases¶
- Re-sequencing after every stage. Extra machinery, same memory, no benefit for stateless stages.
- Filtering without tombstones. The re-sequencer waits forever for the missing number.
- Sequence numbers assigned after a concurrent stage. They record completion order, not input order.
- Unbounded buffers with heavy-tailed latency. Memory grows with the slowest item's duration.
Frequently Asked Questions¶
How do I keep input order through several concurrent asyncio stages?
Tag each item with a sequence number when it enters the pipeline, let stages run concurrently in any order, and re-sequence once with a reorder buffer in front of the consumer that needs order.
Should I reorder after every concurrent stage?
Only before stages whose processing depends on order. In testing, one re-sequencer at the sink held 192 items, the same as two intermediate buffers combined (190 and 10), with the same run time.
What happens if a stage drops an item in an ordered pipeline?
The re-sequencer waits forever for its sequence number. Have every stage emit exactly one envelope per input, using a tombstone for filtered items and the exception for failures.
How do I bound the memory of re-sequencing?
Limit how far the source may run ahead with a semaphore released when items are emitted in order. It caps the buffer but makes stragglers stall the pipeline, so pair it with per-item timeouts.
Related¶
- Async Data Pipelines — up to the topic overview.
- Merging multiple async iterators into one stream — ordering when several sources feed one stream.
- Concurrent Execution & Worker Patterns — the section overview.