Timing Out Each Item of an Async Iterator¶
Streams — a WebSocket feed, a paginated API, rows from a database cursor, a message consumer — should fail when they stall, not when they run long. A timeout around the whole async for loop does the wrong thing: it limits the total duration, so a healthy but long stream is cut off. Tested on Python 3.14 with a stream of items arriving every 50 ms: a 0.5 s timeout per item let all 30 items through in 1.50 s, where a 1.0 s timeout around the loop would have stopped it at item 20. A per-item timeout around anext() caught a 2-second stall at 1.00 s — but afterwards the async generator was finished: the next anext() raised StopAsyncIteration, because the timeout's cancellation had closed it. Pumping the iterator into a queue from a separate task kept the stream alive through the stall: the consumer logged 3 stall warnings at 0.5 s intervals and still received all 16 items. This guide implements both behaviours — fail on stall, or warn and keep waiting.
Prerequisites¶
- Python 3.11+ for
asyncio.timeout; 3.10+ for theanext()builtin. - Timeout tools, from choosing asyncio.timeout vs wait_for.
- Async iterators, from Async Context Managers & Iterators.
1. Do not put the deadline around the whole loop¶
A timeout around async for bounds the stream's total length, not the gap between items:
async with asyncio.timeout(1.0):
async for item in feed(): # a healthy 30-item stream at 50 ms per item
handle(item) # is cut at about 1.0 s, after ~20 items
That is right only when the whole stream genuinely has a deadline — a request that must finish within its budget. For long-lived streams the question is "has it stopped?", which needs a bound on each wait for the next item. Tested: with a per-item bound of 0.5 s, all 30 items of a 1.5-second stream arrived, while a 2-second stall would still have been caught.
Verify: a long, healthy stream in a test runs to completion under the per-item timeout.
2. Time out each item with anext()¶
To fail when the stream stalls, put the timeout around each anext() call:
async def iterate_with_item_timeout(aiter, timeout: float):
iterator = aiter.__aiter__()
while True:
try:
async with asyncio.timeout(timeout):
item = await anext(iterator)
except StopAsyncIteration:
return
yield item
async for message in iterate_with_item_timeout(subscription(), timeout=30):
await handle(message) # TimeoutError if 30 s pass with no message
Tested with a 0.5 s per-item timeout and a 2-second stall after item 10: TimeoutError at 1.00 s, after 10 items. The handler's own time is outside the timeout — only the wait for the next item is bounded — which is usually what you want; if processing time should count too, put the timeout around the handler as well. asyncio.wait_for(anext(iterator), timeout) behaves the same way.
Verify: a test stream that stops yielding raises TimeoutError after one timeout period.
3. Know that the timeout closes the generator¶
When the timeout fires, it cancels the anext() call — and cancelling an async generator while it is suspended inside it finalizes the generator. Tested: after the per-item timeout, the next anext() on the same generator raised StopAsyncIteration; the stream could not be resumed:
iterator = feed()
try:
item = await asyncio.wait_for(anext(iterator), 0.5) # stalls: TimeoutError
except TimeoutError:
pass
await anext(iterator) # tested: StopAsyncIteration - it is over
That is fine when a stall means "give up and reconnect" — then treat the timeout as the end of this stream and start a new one. It is wrong when a stall is merely worth noticing — a market data feed with quiet periods, a log tail — because the first slow period ends the stream. For those, decouple reading from waiting (step 4). Iterators that are not generators (classes with __anext__) behave according to their own code; many network iterators also close their connection on cancellation.
Verify: after a per-item timeout in a test, the code reconnects or restarts the stream rather than calling anext() on the dead iterator.
4. Pump into a queue to warn without closing¶
To observe stalls without killing the stream, let a separate task iterate the source and hand items over through a queue; the consumer's timeout then applies to the queue, not the generator:
_END = object()
async def with_stall_warnings(aiter, timeout: float, on_stall):
queue: asyncio.Queue = asyncio.Queue(maxsize=1)
async def pump():
try:
async for item in aiter:
await queue.put(item)
finally:
await queue.put(_END)
pump_task = asyncio.create_task(pump())
try:
while True:
try:
async with asyncio.timeout(timeout):
item = await queue.get()
except TimeoutError:
on_stall() # the source is untouched and keeps waiting
continue
if item is _END:
break
yield item
await pump_task # re-raise errors from the source
finally:
pump_task.cancel()
Tested with a 2-second stall and a 0.5 s timeout: three stall callbacks at 1.0, 1.5 and 2.0 s, then the remaining items, 16 in total. A queue of size 1 keeps backpressure — the pump cannot read ahead more than one item — so the source still sees a slow consumer as slow. Count consecutive stalls and give up after a limit if an endless silence should still end the stream.
Verify: a test stream with a long pause produces stall warnings and still delivers every item afterwards.
5. Choose per-item, idle and total limits together¶
Long-lived streams often need three bounds: how long to wait for the next item, how long a stream may be idle before it is considered dead, and — for request-scoped streams — a total budget:
async def consume(stream_factory, *, item_timeout=5.0, max_idle=60.0, total=None):
async with asyncio.timeout(total): # None: no total limit
idle = 0.0
async for item in with_stall_warnings(stream_factory(), item_timeout,
on_stall=lambda: None):
idle = 0.0
await handle(item)
A cleaner way to express "idle for too long" is to count consecutive stalls inside on_stall and raise once stalls * item_timeout >= max_idle. The per-item timeout gives early, cheap signals; the idle limit decides when to reconnect; the total budget belongs to request handlers only. Idle connection timeouts at the transport level are covered in implementing idle timeouts for connections.
Verify: each limit is exercised by a test: a short stall (warning only), a long silence (reconnect), and a request budget (fails at the deadline).
Verification¶
Stream timeouts are correct when:
- Long-lived streams bound the gap between items, not their total length.
- Per-item timeouts on generators lead to a reconnect, never to reuse of the closed iterator.
- Streams with quiet periods use a pump task, so stalls warn without closing.
- Idle limits and request budgets are separate, explicit settings.
Diagnostic Hook: log stall warnings with the stream's identity and the time since the last item. A steady trickle of short stalls is normal quiet traffic — raise the per-item timeout; long stalls that end in reconnects on one source point at that upstream; stalls on every stream at once point at your own event loop being blocked.
Pitfalls & edge cases¶
- A timeout around the whole
async for. It cuts healthy long streams. - Calling
anext()after a per-item timeout. Tested:StopAsyncIteration— the generator is closed. - Unbounded pump queues. The pump reads ahead without limit; use
maxsize=1. - Forgetting to re-raise pump errors. Await the pump task at the end.
Frequently Asked Questions¶
How do I add a timeout to each item of an async for loop?
Iterate manually and wrap each anext(iterator) in asyncio.timeout(seconds) or asyncio.wait_for. In testing, a 0.5 s per-item timeout caught a 2-second stall at 1.00 s while letting a long healthy stream complete.
Can I continue an async generator after a timeout on anext()?
No. The timeout cancels the generator while it is suspended, which closes it; in testing the next anext() raised StopAsyncIteration. Reconnect, or use a pump task and queue so the generator is never cancelled.
How do I detect a stalled stream without stopping it?
Let a separate task iterate the stream into a small asyncio.Queue and apply the timeout to queue.get(); a timeout then means only that nothing arrived, and the stream continues.
Should a WebSocket or consumer loop have a total timeout?
Usually not. Long-lived streams need a per-item or idle timeout; total timeouts belong to request-scoped work with a deadline.
Related¶
- Timeouts & Deadlines — up to the topic overview.
- Telling TimeoutError apart from CancelledError — what the per-item timeout raises, and why.
- Resilience, Cancellation & Error Handling — the section overview.