Skip to content

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

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.

Run time by prefetch depth A grid of 4 rows by 6 columns. Run time by prefetch depth source / consumer depth 0 depth 1 depth 2 depth 4 depth 8 50 ms pages / 40 ms processing, 40 pages 3.72 s 2.09 s 2.09 s 2.09 s - 50 ms pages / 20 ms processing, 40 pages 2.87 s 2.05 s 2.06 s 2.06 s - 20 ms pages / 40 ms processing, 40 pages 2.50 s 1.73 s 1.70 s 1.73 s - jittery 10-170 ms pages / 40 ms, 80 pages 7.96 s 5.81 s 5.40 s 4.92 s 4.58 s Python 3.14; latencies simulated with asyncio.sleep.

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.

Prefetch with an early exit A sequence of 7 messages between 4 participants. Prefetch with an early exit consumer queue (depth 4) pump task source fetch page 1 put page 1 get page 1, process fetch page 2 (overlaps) break after page 3 context exit: cancel pump aclosing: source closed The context manager guarantees the pump and source end with the block.

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.

How much prefetch does this source need? A decision on What are the two sides like with 4 outcomes. How much prefetch does this source need? What are the two sides like? consumer does no I/O no prefetch nothing to overlap steady latency depth 1 3.72 s to 2.09 s bursty / jittery source depth 4-8 7.96 s to 4.58 s large items smallest depth that helps memory = depth x item Always as a context manager, so early exits clean up.

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 BaseException in 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.