Adapting Blocking Iterators to Async Iterators¶
Plenty of iterators block: a synchronous database cursor, csv.reader over a network file system, a cloud SDK's paginator that makes an HTTP call every hundred items, a generator that reads from a serial port. Dropped into a coroutine as for row in cursor:, each next() that blocks stalls the event loop. The obvious fix — await asyncio.to_thread(next, it) for every item — works and is slow: in a test, iterating 20,000 items that way cost 43.3 µs per item, almost all of it thread hand-off overhead. Pulling items in batches of 500 per thread hop cost 0.11 µs per item. This guide builds an adapter that batches, then one that prefetches on a dedicated thread, and shows how to close the underlying iterator when the consumer stops early.
Prerequisites¶
- Python 3.11+, stdlib only.
- Thread offloading, from running blocking SDK calls with asyncio.to_thread.
- Async iterator protocol, from building async iterator classes.
1. Start with the per-item adapter, and measure it¶
The simplest correct adapter moves each next() call into the default executor:
import asyncio
from collections.abc import AsyncIterator, Iterator
from typing import TypeVar
T = TypeVar("T")
_DONE = object()
async def iterate_in_thread(it: Iterator[T]) -> AsyncIterator[T]:
while (item := await asyncio.to_thread(next, it, _DONE)) is not _DONE:
yield item
The sentinel default matters: next(it) raising StopIteration inside a thread would surface through the future as a StopIteration, which Python converts to RuntimeError when it propagates out of a coroutine. Passing a default makes exhaustion an ordinary return value.
It is correct, and every item pays for two thread hand-offs plus a future. For iterators whose items are each expensive to produce — one network round trip per item — that overhead is irrelevant. For iterators that produce thousands of cheap items between occasional slow ones, such as a cursor that fetches a page then yields rows from memory, the overhead dominates.
Verify: time 10,000 items through the adapter against a plain for loop over the same iterator; the difference is the per-item thread overhead.
2. Batch the thread hops¶
Pull a batch of items per thread call with itertools.islice, and yield them from memory on the loop:
import itertools
async def iterate_batched(it: Iterator[T], batch_size: int = 500) -> AsyncIterator[T]:
def take() -> list[T]:
return list(itertools.islice(it, batch_size))
while batch := await asyncio.to_thread(take):
for item in batch:
yield item
Measured on the same 20,000 items: 0.11 µs per item, about 400 times cheaper. The trade-offs are memory — up to batch_size items are held at once — and latency to the first item, which now waits for a whole batch. For a cursor over a slow query, that first-batch wait is the time to fetch 500 rows rather than one.
Choose the batch size so a batch takes a few milliseconds to produce: large enough to amortise the hop, small enough that the consumer sees steady progress. If the underlying iterator already has a natural page size — a cursor's arraysize, an SDK's page — match it.
Verify: log the time each take() call spends in the thread; aim for single-digit milliseconds per batch.
3. Prefetch on a dedicated thread¶
Batching still alternates: the loop waits while a batch is produced, then the thread idles while the consumer processes it. When the consumer does real async work per item — writing each row to another service — a producer thread that keeps reading ahead overlaps the two:
import queue
import threading
async def prefetch(it: Iterator[T], maxsize: int = 1000) -> AsyncIterator[T]:
loop = asyncio.get_running_loop()
q: asyncio.Queue = asyncio.Queue(maxsize=maxsize)
stop = threading.Event()
def producer() -> None:
try:
for item in it:
if stop.is_set():
return
asyncio.run_coroutine_threadsafe(q.put(item), loop).result() # blocks when full
asyncio.run_coroutine_threadsafe(q.put(_DONE), loop).result()
except BaseException as exc: # hand errors over
asyncio.run_coroutine_threadsafe(q.put(exc), loop).result()
thread = threading.Thread(target=producer, name="prefetch", daemon=True)
thread.start()
try:
while (item := await q.get()) is not _DONE:
if isinstance(item, BaseException):
raise item
yield item
finally:
stop.set()
while not q.empty(): # unblock a producer waiting on a full queue
q.get_nowait()
The bounded queue gives backpressure: when the consumer falls behind, q.put blocks the producer thread instead of letting it read the whole source into memory. Errors raised in the thread are passed through the queue and re-raised in the consumer, so a failed read is not silently treated as end-of-data. The bridging details are covered in bridging queues between threads and asyncio tasks.
Verify: with a slow consumer, memory stays bounded by maxsize items; with a failing source, the consumer receives the source's exception.
4. Close the source when the consumer stops early¶
A consumer that breaks out of the loop leaves the source open: the cursor still holds its server-side resources, the file is still open, the producer thread is still blocked on q.put. The finally in each adapter only runs when the async generator itself is closed — which, after a break, is whenever the generator is finalised, unless the consumer uses aclosing():
from contextlib import aclosing
async def first_matching(cursor, predicate):
async with aclosing(prefetch(iter(cursor))) as rows:
async for row in rows:
if predicate(row):
return row # aclosing runs the finally now
The adapter's finally should also close the underlying iterator if it supports it. Many blocking iterators expose close() (generators, DB-API cursors, file objects); call it in the thread, since closing may itself block:
async def prefetch(it, maxsize: int = 1000):
... # as above
try:
...
finally:
stop.set()
close = getattr(it, "close", None)
if close is not None:
await asyncio.to_thread(close)
For the prefetching adapter, the queue must also be drained after setting stop, as in the code above, or a producer blocked on a full q.put never observes the flag. Verified: breaking out after the fourth of a million items left no producer thread running. The general rule for generator cleanup is in closing async generators with aclosing.
Verify: break out after the first item; the source's close() is called before the function returns and the producer thread exits.
5. Pick the adapter by the source's cost profile¶
The three adapters fit different sources:
| Source | Per-item cost | Adapter |
|---|---|---|
one network call per item (SDK get per id) |
milliseconds | per-item to_thread |
| cheap items with occasional slow fetches (cursor, paginator) | microseconds, with spikes | batched to_thread |
| steady slow stream with slow consumer work (file over NFS, serial port) | sustained | prefetching thread |
| an async-native alternative exists | — | use it instead |
The last row is the most important. Async drivers exist for most databases, and many SDKs have async clients; an adapter is the bridge while you cannot use them, as in the asyncpg cursor approach of streaming large result sets with asyncpg cursors. For file reads specifically, compare with async file I/O with aiofiles vs asyncio.to_thread.
Verify: profile the adapter under realistic load; thread hand-off time should be a small fraction of total iteration time.
Verification¶
The adapter is right when:
- The event loop never blocks on the source — a lag monitor beside it stays flat.
- Per-item overhead is small relative to per-item work.
- Memory is bounded by batch size or queue size, not by source length.
- Early termination closes the source and stops any producer thread.
- Source errors reach the consumer as exceptions, never as a silent end of iteration.
Diagnostic Hook: export the queue depth of a prefetching adapter and the time per batch of a batched one. A prefetch queue that sits at maxsize means the consumer is the bottleneck; one that sits at zero means the source is. Batch times that suddenly grow usually mean the default executor is saturated — check queue wait as in sizing the default thread pool executor.
Pitfalls & edge cases¶
- Calling
next(it)without a default in a thread.StopIterationcrossing a coroutine boundary becomesRuntimeError. - Using one blocking iterator from several threads. Most are not thread-safe; the batched and prefetching adapters each touch it from one thread at a time, keep it that way.
- Iterators bound to their creating thread. Some DB drivers and SQLite connections refuse use from a different thread; create them inside the producer thread.
- Unbounded prefetch.
asyncio.Queue()withoutmaxsizeturns a slow consumer into a memory leak.
Frequently Asked Questions¶
How do I iterate a blocking iterator from asyncio?
Run the blocking next() calls in a thread. For cheap items, pull them in batches with itertools.islice inside asyncio.to_thread and yield from the batch; for one-at-a-time expensive items, call to_thread(next, it, sentinel) per item.
Why is to_thread per item so slow?
Each call hands work to a thread pool and waits for a future to resolve, which cost about 43 µs per item in testing — far more than producing a cheap item. Batching 500 items per call brought it to 0.11 µs per item.
How do I stop a prefetching thread when the consumer breaks out?
Wrap the async generator in contextlib.aclosing so its finally block runs immediately, set a stop flag there, drain the queue so a blocked producer can see the flag, and close the underlying iterator in a thread.
Why use a sentinel with next() in to_thread?
Because StopIteration raised inside a coroutine is converted to RuntimeError. next(it, sentinel) returns the sentinel at exhaustion, which the adapter checks to end the loop normally.
Related¶
- Async Context Managers & Iterators — up to the topic overview.
- Async generators vs queues for streaming pipelines — the choice the prefetching adapter makes internally.
- Asyncio Fundamentals & Event Loop Architecture — the section overview.