Prefetching Items from Async Iterators¶
An async for over a paginated API or a database cursor alternates two waits: fetch a page, process it, fetch the next. The fetch for page n+1 does not start until page n is processed, so the total is the sum of both. Prefetching runs the source in a background task that stays a few items ahead, so fetching and processing overlap. Measured on Python 3.14 with a source that took 50 ms per page and a consumer that spent about 40 ms processing each page: 40 pages took 3.72 s without prefetch and 2.09 s with one page prefetched — close to the 2.0 s the fetches alone take. Deeper prefetch made no difference with constant latency: 2.09 s at depths 2 and 4. With a jittery source — three pages at 10 ms, then one at 170 ms — depth helped: 80 pages took 7.96 s unbuffered, 5.81 s at depth 1, 4.92 s at depth 4 and 4.58 s at depth 8. A correctly built prefetcher also closed its source and left no tasks behind when the consumer stopped after three pages, and passed a source error through to the consumer. This guide builds that prefetcher.
Prerequisites¶
- Python 3.11+, an async iterator that waits between items.
- Closing generators, from closing async generators with aclosing.
- The topic overview, Async Context Managers & Iterators.
1. See the serial wait in a plain async for¶
A paginated source and a consumer that does async work per item:
async def pages(client, url):
while url:
resp = await client.get(url) # ~50 ms per page
body = resp.json()
yield body["items"]
url = body.get("next")
async def sync_all(client, url):
async for items in pages(client, url):
for item in items:
await store(item) # ~4 ms each, ~40 ms per page
The generator is suspended at yield while the consumer stores items, so the request for the next page is not even sent until storing finishes. Measured with simulated 50 ms pages and 40 ms of processing: 3.72 s for 40 pages, about 93 ms each — the sum of the two. Whenever both sides wait on I/O, that sum can become a maximum instead.
Verify: time per item without prefetch is close to fetch time plus processing time.
2. Run the source in a background task with a bounded queue¶
Move iteration of the source into a task that puts items on a bounded queue, and have the consumer read from the queue. The queue's size is the prefetch depth — how far the producer may run ahead:
import asyncio
from contextlib import aclosing, asynccontextmanager
_DONE = object()
@asynccontextmanager
async def prefetch(source, depth: int = 1):
queue: asyncio.Queue = asyncio.Queue(maxsize=depth)
async def pump():
try:
async with aclosing(source) as src:
async for item in src:
await queue.put((item, None)) # waits when depth items are ready
await queue.put((_DONE, None))
except Exception as exc:
await queue.put((_DONE, exc)) # hand the error to the consumer
task = asyncio.create_task(pump())
async def items():
while True:
item, exc = await queue.get()
if exc is not None:
raise exc
if item is _DONE:
return
yield item
try:
yield items()
finally:
task.cancel() # consumer finished or left early
try:
await task
except asyncio.CancelledError:
pass
Used as:
async with prefetch(pages(client, url), depth=1) as it:
async for items in it:
for item in items:
await store(item)
Measured: 2.09 s for the 40 pages, against 3.72 s without prefetch. The producer fetches page n+1 while the consumer stores page n, and the queue bound stops it from fetching the whole API into memory when the consumer is slow.
Verify: with prefetch, total time approaches the larger of total fetch time and total processing time.
3. Choose the depth from the source's variability¶
With steady latency, one page of prefetch captures the whole gain: measured, depths 1, 2 and 4 all took 2.05–2.09 s when the source was the slower side, and 1.70–1.74 s when the consumer was. More depth only buffers items that the slower side cannot use yet. With variable latency, depth absorbs the variation: while the producer waits on a slow page, the consumer drains items fetched during fast ones.
DEPTH = 4 # covers bursts of slow pages; costs up to 4 pages of memory
Measured with a source that returned three pages in 10 ms and then one in 170 ms: 5.81 s at depth 1, 5.40 s at depth 2, 4.92 s at depth 4 and 4.58 s at depth 8, against a lower bound of a little over 4 s set by the total fetch time. Memory grows linearly with depth — depth × page size — so pick the smallest depth that covers the typical run of slow pages, and measure on real latency data, since simulated jitter only approximates it.
Verify: increasing depth further no longer reduces run time on representative latency.
4. Close the source and the task when the consumer stops early¶
A prefetcher creates a task, and the consumer may stop at any time — break, an exception, cancellation. The background task must stop, and the source must be closed, or a task keeps fetching pages nobody reads:
async with prefetch(pages(client, url), depth=4) as it:
async for items in it:
if found(items):
break # the context manager's exit cancels the pump
Measured with depth 4 and a consumer that broke after three pages: the source had fetched 3 pages, its finally block had run, and no tasks remained. This is why prefetch is a context manager rather than a plain async generator: the finally around yield items() runs when the async with block ends, in the consumer's task, and cancels the pump there. A prefetching async generator that started its own task would instead rely on the generator being closed — which, without aclosing, happens late and in a different task. The aclosing(source) inside the pump makes sure that cancelling the pump also closes the source promptly, releasing cursors and connections.
Verify: after a consumer breaks out early, no pump task remains and the source's cleanup has run.
5. Pass source errors to the consumer¶
If the source fails, the consumer must see the error, not hang waiting for an item that will never come. The pump catches the exception and puts it on the queue in place of an item; the consumer re-raises it:
try:
async with prefetch(pages(client, url), depth=2) as it:
async for items in it:
await handle(items)
except ConnectionError as exc:
log.warning("sync stopped: %s", exc)
Measured with a source that raised ConnectionError on its sixth page: the consumer received "page 5 failed" after handling the pages before it, and the source's cleanup had run. Items already in the queue when the error happened are delivered first, so the consumer sees the error in order. The pump catches Exception, not BaseException, so cancellation of the pump itself is not mistaken for a source error. For consumers that need several sources at once, the same queue-and-pump structure generalises to merging, as in merging multiple async iterators into one stream.
Verify: an exception in the source reaches the consumer, and nothing is left running.
Verification¶
Prefetching is correct and worthwhile when:
- Both sides wait on I/O, and run time with prefetch approaches the slower side's total.
- Depth is the smallest that removes the gain on representative latency.
- The prefetcher is a context manager that cancels its task and closes the source on exit.
- Source errors surface in the consumer, in order.
Diagnostic Hook: sample the prefetch queue's size. A queue that is always full means the consumer is the bottleneck and extra depth only costs memory; one that is always empty means the source is the bottleneck and the consumer waits on every item — prefetch depth then cannot help, but fetching pages concurrently, if the API allows, can.
Pitfalls & edge cases¶
- Expecting depth to fix a slow source. Measured: depth 1, 2 and 4 all took 2.09 s.
- A prefetching async generator without a context manager. Its task outlives an early exit.
- Catching
BaseExceptionin the pump. Cancellation would be reported as a source error. - Unbounded queues. The pump fetches everything if the consumer is slow.
Frequently Asked Questions¶
How do I prefetch items from an async iterator?
Iterate the source in a background task that puts items on an asyncio.Queue with maxsize equal to the prefetch depth, and consume from the queue. One page of prefetch cut 40 pages from 3.72 s to 2.09 s in testing.
How deep should prefetching be?
One item captures the gain when latency is steady; depths 2 and 4 added nothing. With jittery latency, deeper buffers helped: 5.81 s at depth 1, 4.58 s at depth 8.
How do I stop a prefetch task when the consumer breaks out?
Wrap the prefetcher in an async context manager whose exit cancels the task and closes the source with aclosing; after a break, no task remained and the source was closed.
What happens if the source raises during prefetch?
The pump puts the exception on the queue and the consumer re-raises it after the items fetched before the failure.
Related¶
- Async Context Managers & Iterators — up to the topic overview.
- Writing reentrant async context managers — context managers that nest.
- Asyncio Fundamentals & Event Loop Architecture — the section overview.