Skip to content

Writing Async itertools Helpers: take, chunked and a Bounded amap

itertools has no async counterpart in the standard library, so every codebase that streams data through async generators ends up writing the same handful of helpers: take the first N items, group items into batches, map a coroutine over a stream with bounded concurrency. They look like five-line functions. Written naively, each has a subtle failure — take leaves its source open, a time-based chunked loses an item when its timer fires, a concurrent amap either runs unbounded or returns results out of order. This guide writes the three that matter most, each verified: take closed its 100-item source after yielding 3; chunked flushed a partial batch 100 ms into a pause; and an ordered amap with a limit of 10 processed 100 items in 0.41 s against 2.91 s sequentially, peaking at exactly 10 in flight with results in input order.

Prerequisites

1. Write take() so it closes its source

take(src, n) yields the first n items. The part that matters is what happens to the source after the n-th item:

from collections.abc import AsyncIterable, AsyncIterator
from contextlib import aclosing
from typing import TypeVar

T = TypeVar("T")


async def take(src: AsyncIterator[T], n: int) -> AsyncIterator[T]:
    if n <= 0:
        return
    async with aclosing(src) as it:
        count = 0
        async for item in it:
            yield item
            count += 1
            if count >= n:
                return                            # aclosing closes src right here

Without aclosing, stopping after n items leaves the source generator suspended at its last yield, holding whatever it holds — a database cursor, an HTTP stream, a subscription — until the garbage collector finalises it. With it, the source's finally runs as soon as take returns. Verified: taking 3 items from a 100-item source ran the source's finally before the next statement.

Note that aclosing requires an object with aclose(), which async generators have and many iterator classes do not. For a general AsyncIterable, close only when the method exists.

Verify: give take a source with a finally that appends to a list; the list is updated before take's consumer moves on.

take(src, 3) closes its source after the third item A sequence of 5 messages between 3 participants. take(src, 3) closes its source after the third item consumer take source next item anext() items 0, 1, 2 count reached 3: return aclose(): finally runs now Closing the source is part of take's job, not something to leave to garbage collection.

2. Write chunked() by size and by time

Batching by size alone stalls when the stream pauses: four items wait forever for a fifth. Production batching flushes on whichever comes first, a full batch or a maximum wait since the batch's first item. The trap is implementing the time limit with asyncio.wait_for(anext(it), timeout) — when the timeout fires, it cancels the pending anext, which can lose the item the source was about to produce. Keep the pending anext alive across timeouts instead:

import asyncio


async def chunked(src: AsyncIterable[T], size: int, max_wait: float) -> AsyncIterator[list[T]]:
    it = aiter(src)
    loop = asyncio.get_running_loop()
    buf: list[T] = []
    pending: asyncio.Future | None = None
    deadline = 0.0
    try:
        while True:
            if pending is None:
                pending = asyncio.ensure_future(anext(it))
            timeout = max(0.0, deadline - loop.time()) if buf else None
            done, _ = await asyncio.wait({pending}, timeout=timeout)
            if not done:                          # time is up: flush, keep waiting on the same anext
                yield buf
                buf = []
                continue
            fut, pending = pending, None
            try:
                item = fut.result()
            except StopAsyncIteration:
                if buf:
                    yield buf
                return
            if not buf:
                deadline = loop.time() + max_wait
            buf.append(item)
            if len(buf) >= size:
                yield buf
                buf = []
    finally:
        if pending is not None:
            pending.cancel()

Verified on a stream that produced 7 items, paused 300 ms, then produced 2 more, with size=5, max_wait=0.1: batches [0..4] at 0 ms, [5, 6] at 100 ms, and [7, 8] at the end of the stream. The same logic over a queue rather than an iterator is in batching queue items by size and time.

Verify: no item is lost or duplicated across timeouts — compare the concatenated batches with the source.

chunked(size=5, max_wait=100 ms) on a bursty stream 2 lanes over time. chunked(size=5, max_wait=100 ms) on a bursty stream source items 0-6 300 ms pause 7, 8 batches [0..4] full [5, 6] by timer [7, 8] end time → Size and time limits together bound both batch size and how long an item can wait in a batch.

3. Write an ordered amap with bounded concurrency

amap(fn, src, limit) applies a coroutine function to each item with at most limit calls in flight, yielding results in input order. A sliding window of tasks does all three:

from collections import deque
from collections.abc import Awaitable, Callable

R = TypeVar("R")


async def amap(fn: Callable[[T], Awaitable[R]], src: AsyncIterable[T], limit: int) -> AsyncIterator[R]:
    window: deque[asyncio.Task[R]] = deque()
    try:
        async for item in src:
            window.append(asyncio.create_task(fn(item)))
            if len(window) >= limit:
                yield await window.popleft()       # wait for the OLDEST, preserving order
        while window:
            yield await window.popleft()
    finally:
        for t in window:
            t.cancel()
        await asyncio.gather(*window, return_exceptions=True)

The window never holds more than limit tasks, so concurrency is bounded without a semaphore, and pulling from the source pauses while the window is full — backpressure on the source for free. Results come out in input order because the generator always awaits the oldest task. The cost of ordering is head-of-line blocking: one slow item holds back faster ones behind it, which is why the measured run took 0.41 s rather than the theoretical minimum for 100 items of 10–50 ms at concurrency 10.

The finally makes early exit safe: if the consumer stops, or one call raises, every in-flight task is cancelled and awaited. Pair it with aclosing() at the call site so that happens immediately.

Verify: count in-flight calls inside fn; the peak equals limit, and output order equals input order.

100 calls of 10-50 ms each 2 horizontal bars comparing sequential await with the others. 100 calls of 10-50 ms each sequential await 2.91 s ordered amap, limit 10 0.41 s, peak 10 in flight Random per-call latency between 10 and 50 ms; output order verified equal to input order. Bounded concurrency gives most of the speed-up; strict ordering costs some head-of-line waiting.

4. Offer an unordered variant when order does not matter

When results are independent, drop the ordering and yield each result as it finishes. It removes head-of-line blocking at the cost of losing the input mapping, so yield pairs:

async def amap_unordered(fn, src: AsyncIterable[T], limit: int):
    running: set[asyncio.Task] = set()
    it = aiter(src)
    exhausted = False
    try:
        while True:
            while not exhausted and len(running) < limit:
                try:
                    item = await anext(it)
                except StopAsyncIteration:
                    exhausted = True
                    break
                task = asyncio.create_task(fn(item))
                task.item = item                    # keep the input next to its result
                running.add(task)
            if not running:
                return
            done, running = await asyncio.wait(running, return_when=asyncio.FIRST_COMPLETED)
            for t in done:
                yield t.item, t.result()           # re-raises fn's exception here
    finally:
        for t in running:
            t.cancel()
        await asyncio.gather(*running, return_exceptions=True)

This is the same loop-on-wait(FIRST_COMPLETED) shape described in using asyncio.wait with FIRST_COMPLETED, packaged as an iterator. Setting an attribute on a task is allowed and is the cheapest way to carry the input along.

Verify: with skewed latencies, the unordered variant finishes no later than the ordered one, and every input appears exactly once in the output.

5. Decide how errors propagate

Both amap variants re-raise the first failing call's exception at the point its result is yielded, and the finally then cancels the rest — fail-fast semantics, matching a TaskGroup. Pipelines that must continue past bad items need the error as data instead:

from dataclasses import dataclass


@dataclass
class Failed:
    item: object
    error: BaseException


def capture(fn):
    async def wrapper(item):
        try:
            return await fn(item)
        except Exception as exc:                  # never CancelledError
            return Failed(item, exc)
    return wrapper


async for result in amap(capture(enrich), records, limit=20):
    if isinstance(result, Failed):
        await dead_letters.put(result)            # handle, do not stop the stream
        continue
    await sink.write(result)

Catching Exception rather than BaseException keeps cancellation working. Failed items can then go to a dead-letter path as in implementing a dead letter queue with asyncio.

Verify: inject one failing item; with capture the stream completes and the failure is recorded, without it the stream stops and no task is left running.

Verification

The helpers are correct when:

  • Sources are closed promptly when a helper stops early.
  • No items are lost or duplicated, including across chunked timeouts.
  • amap concurrency never exceeds limit, and ordered output matches input order.
  • Early exit and errors leave zero tasks behind in asyncio.all_tasks().

Diagnostic Hook: in pipelines built from these helpers, record per-stage counters — items in, items out, batches flushed by size versus by time, and in-flight calls in amap. A chunked stage that almost always flushes by time is undersized for its traffic; an amap stage whose in-flight count sits at limit is the pipeline's bottleneck.

Pitfalls & edge cases

  • asyncio.wait_for(anext(it), t) for time-based batching. Cancelling anext on timeout can drop an item mid-production.
  • Unbounded create_task per item. A fast source creates tasks faster than they finish; always bound with a window or limit.
  • Iterating a source in two helpers at once. Async generators raise RuntimeError: anext(): asynchronous generator is already running when two consumers pull concurrently.
  • Forgetting aclosing at the call site. The helper's finally then waits for garbage collection like any other generator's.

Frequently Asked Questions

Is there an async version of itertools in the standard library?

No. The standard library offers aiter, anext and async comprehensions, but no async itertools. Third-party packages such as aiostream and asyncstdlib provide them, or you can write the few you need with explicit cleanup.

How do I batch an async stream by size and time?

Keep one pending anext future across timeouts, wait on it with asyncio.wait and a timeout measured from the batch's first item, flush on timeout or when the batch is full, and never cancel the pending anext except at shutdown.

How do I map a coroutine over an async iterator with limited concurrency?

Keep a sliding window of at most limit tasks. To preserve order, await the oldest task before adding a new one; for unordered results, loop on asyncio.wait with FIRST_COMPLETED. Cancel and await the window in a finally block.

Why does my async generator stay open after I take a few items?

Breaking out of async for leaves the generator suspended until it is finalised. Wrap it in contextlib.aclosing so its cleanup runs as soon as you stop.