Merging Multiple Async Iterators into One Stream¶
Pipelines often have several sources that a single consumer should treat as one stream: pages from several API cursors, messages from several queues, lines from several log files, events from several websockets. Reading them one after another wastes the concurrency asyncio offers — the consumer waits on a slow source while a fast one has data ready. A merge reads all of them concurrently and yields whatever arrives first. Getting the lifecycle right is the real work: a merge that does not close its sources when the consumer stops early leaks connections, and one that swallows a failing source's exception silently produces partial output. Tested with a fast source (30 items at 10 ms) and a slow one (3 items at 100 ms), the merge yielded all 33 items in 0.31 s — the slow source's duration, not the sum; breaking out after five items closed both sources; and a source that raised ConnectionError mid-stream had its error surfaced to the consumer. This guide builds that merge.
Prerequisites¶
- Python 3.11+, stdlib only.
- Generator cleanup, from closing async generators with aclosing.
- TaskGroup error semantics, from handling ExceptionGroup from TaskGroup.
1. Pump each source into a shared bounded queue¶
One task per source pulls items and puts them into a shared queue; the merge yields from the queue. A sentinel marks the end:
import asyncio
from contextlib import aclosing
async def merge(*sources, maxsize: int = 100):
q: asyncio.Queue = asyncio.Queue(maxsize)
DONE = object()
async def pump(src):
async with aclosing(src) as it: # close the source however the pump ends
async for item in it:
await q.put(item) # blocks when the consumer is behind
async def run_all():
try:
async with asyncio.TaskGroup() as tg:
for s in sources:
tg.create_task(pump(s))
finally:
await q.put(DONE) # always tell the consumer we are done
runner = asyncio.create_task(run_all())
try:
while (item := await q.get()) is not DONE:
yield item
await runner # re-raise any source failure
finally:
runner.cancel()
await asyncio.gather(runner, return_exceptions=True)
The bounded queue gives backpressure to every source at once: when the consumer is slow, all pumps block on put, so no source runs ahead into memory. Measured: 30 fast items and 3 slow ones merged in 0.31 s, the duration of the slow source alone.
Verify: merge two sources with different rates; total time equals the slowest source's duration, and every item arrives.
2. Close every source when the consumer stops early¶
A consumer that breaks out of the loop — found what it needed, hit a limit — must not leave the pumps running. The merge's own finally cancels the runner, which cancels the pumps, and each pump's aclosing closes its source. But that finally only runs when the merge generator itself is closed, so the consumer must close it promptly:
from contextlib import aclosing
async def first_errors(sources, limit: int = 5) -> list:
found = []
async with aclosing(merge(*sources)) as stream:
async for event in stream:
if event.level == "error":
found.append(event)
if len(found) >= limit:
break # aclosing -> merge finally -> pumps cancelled
return found
Verified: breaking out after five items with aclosing ran both sources' finally blocks. Without aclosing, the merge generator would be closed only when garbage collected, and until then two pump tasks and their sources — open HTTP streams, subscriptions — would stay alive.
Verify: give each source a finally that records its closing; break early and check that every source is recorded as closed before the function returns.
3. Surface source failures¶
When one source fails, the consumer must find out. In this merge, a failing pump raises inside the TaskGroup, which cancels the other pumps; the runner's finally still puts DONE, the consumer's loop ends, and await runner re-raises the group's exception:
async def consume(sources):
try:
async for item in merge(*sources):
await handle(item)
except* ConnectionError as eg:
log.error("a source failed: %s", eg.exceptions[0])
raise
Verified: a source that raised ConnectionError("source b failed") after one item surfaced as an ExceptionGroup containing that error, after the items already queued had been yielded. The other source was cancelled. If you would rather keep the healthy sources running and only report the failed one, catch exceptions inside pump and put an error marker into the queue instead of letting the TaskGroup fail.
Verify: inject a failure in one source; the consumer receives the exception and no pump task remains afterwards.
4. Tag items with their source when it matters¶
Arrival order loses the information about which source an item came from. When the consumer needs it — to acknowledge a message to the right queue, to checkpoint per source — tag items in the pump:
async def merge_tagged(named: dict[str, object], maxsize: int = 100):
async def tagged(name, src):
async with aclosing(src) as it:
async for item in it:
yield name, item
async for pair in merge(*(tagged(n, s) for n, s in named.items()), maxsize=maxsize):
yield pair
async for source_name, msg in merge_tagged({"orders": orders_stream, "refunds": refunds_stream}):
await process(msg)
await ack(source_name, msg)
Per-source checkpoints follow naturally: record the last processed position per source_name, so a restart resumes each source independently, as in checkpointing progress in long-running async jobs.
Verify: every item's tag matches its origin, and acknowledgements go to the right source.
5. Merge in order when sources are individually sorted¶
When each source is sorted — log files by timestamp, event streams by sequence — and the output must be globally sorted, arrival order is wrong. A k-way merge reads the next item from each source and always yields the smallest:
import heapq
async def merge_sorted(*sources, key=lambda x: x):
iters = [aiter(s) for s in sources]
heap = []
for idx, it in enumerate(iters):
try:
first = await anext(it)
heapq.heappush(heap, (key(first), idx, first))
except StopAsyncIteration:
pass
while heap:
_, idx, item = heapq.heappop(heap)
yield item
try:
nxt = await anext(iters[idx])
heapq.heappush(heap, (key(nxt), idx, nxt))
except StopAsyncIteration:
pass
The ordered merge holds exactly one item per source, so memory is tiny, but it cannot yield until it has a candidate from every live source: the slowest source paces the whole stream. That is the price of global order. The source index in each heap entry is a tie-breaker so equal keys never compare the items themselves. The same head-of-line trade-off appears in preserving order across concurrent pipeline stages.
Verify: merging sorted sources yields a sorted stream, and the stream waits on the slowest source.
Verification¶
The merge is correct when:
- All items from all sources arrive, in the order the merge promises.
- Early exit closes every source, verified by their
finallyblocks. - A failing source's exception reaches the consumer, and no pump outlives the merge.
- The queue is bounded, so a slow consumer slows every source.
Diagnostic Hook: export items per second per source, tagged at the pump, and the shared queue's depth. A source whose rate drops to zero while others continue is stalled or disconnected; a queue at its bound means the consumer is the bottleneck, and all sources are being throttled to its pace.
Pitfalls & edge cases¶
- No
aclosingat the call site. Early exit leaves pumps and sources running until garbage collection. - Unbounded shared queue. A fast source fills memory while the consumer works.
- Swallowing source exceptions. The output silently lacks one source's data.
- Heap entries without a tie-breaker. Equal keys compare the items and may raise
TypeError.
Frequently Asked Questions¶
How do I merge several async iterators in Python?
Start one task per source that puts items into a shared bounded asyncio.Queue, yield from the queue until every source has finished, and cancel the tasks in a finally block. In testing, merging a 0.3 s source and a 0.31 s source took 0.31 s.
Does breaking out of a merged async stream close the sources?
Only if the merge cancels its pump tasks and the consumer closes the merge promptly. Use contextlib.aclosing around the merged stream so its cleanup runs immediately; each pump should close its own source with aclosing too.
What happens if one source fails during a merge?
In a TaskGroup-based merge, the other sources are cancelled and the failure is re-raised to the consumer as an ExceptionGroup after already-queued items are yielded. Catch errors inside the pump instead if healthy sources should keep running.
How do I merge sorted async streams into one sorted stream?
Use a k-way merge with a heap holding one item per source, always yielding the smallest and refilling from its source. It needs very little memory but moves at the pace of the slowest source.
Related¶
- Async Data Pipelines — up to the topic overview.
- Fanning out a queue to multiple consumer groups — the opposite direction.
- Concurrent Execution & Worker Patterns — the section overview.