Skip to content

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

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.

Timeout placement on a stream, tested A grid of 3 rows by 3 columns. Timeout placement on a stream, tested approach stalled stream (2 s gap at item 10) healthy 30-item stream timeout around whole loop (1.0 s) stopped at 1.00 s, 10 items cut at ~20 items timeout around anext() (0.5 s) stopped at 1.00 s; generator closed all 30 items, 1.50 s pump task + queue (0.5 s) 3 warnings, all 16 items all items Python 3.14; items every 50 ms.

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.

Keeping a stream alive through stalls A flow of 5 stages. Keeping a stream alive through stalls pump task async for over source queue(maxsize=1) hands items over consumer timeout on queue.get() timeout -> warn source untouched stream resumes items flow again Tested: three warnings during a 2 s stall, then all 16 items arrived.

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).

Which timeout should this stream have? A decision on What should a stall mean with 4 outcomes. Which timeout should this stream have? What should a stall mean? request deadline timeout around the loop total duration stall = reconnect timeout around anext() generator then closed stall = just worth noting pump task + queue timeout stream survives long-lived warnings + idle limit reconnect after N stalls Bound the gap between items, not the life of the stream.

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.