Chunking Large Inputs for Async Batch Processing¶
Most sinks a pipeline writes to — databases, object stores, search indexes, bulk APIs — charge a fixed cost per call: a network round trip, a statement parse, a transaction commit, an HTTP request. Writing items one at a time pays that cost per item; writing chunks pays it per chunk. The difference is routinely two orders of magnitude. Measured with asyncpg against PostgreSQL 17, inserting rows one per call managed 8,648 rows per second; executemany with chunks of 100 managed 196,871, of 1,000 306,279 and of 10,000 400,527; COPY with chunks of 10,000 managed 1,543,991, and COPY chunks of 5,000 spread across four connections 2,833,049. This guide covers chunking the input, choosing the chunk size, bounding chunks by bytes as well as count, and running chunks concurrently without overwhelming the sink.
Prerequisites¶
- Python 3.11+; the database examples use
pip install asyncpgand PostgreSQL. - Pipeline stages, from building a staged async pipeline with bounded queues.
- Bulk loading specifics, from bulk loading rows with asyncpg COPY.
1. Chunk an iterable without materialising it¶
Chunking must work on inputs larger than memory — a cursor over millions of rows, a multi-gigabyte file — so it must be lazy:
import itertools
from collections.abc import AsyncIterable, AsyncIterator, Iterable, Iterator
def chunks(iterable: Iterable, size: int) -> Iterator[list]:
it = iter(iterable)
while batch := list(itertools.islice(it, size)):
yield batch
async def achunks(source: AsyncIterable, size: int) -> AsyncIterator[list]:
batch: list = []
async for item in source:
batch.append(item)
if len(batch) >= size:
yield batch
batch = []
if batch:
yield batch # the final partial chunk
Python 3.12 added itertools.batched(iterable, n), which yields tuples and covers the synchronous case directly. The async version must remember the final partial chunk — a bug that silently drops the last few hundred records is easy to write and hard to notice. For streams that pause, add a time limit so a partial chunk does not wait forever, as in writing async itertools helpers.
Verify: for an input of 10,001 items and chunk size 1,000, eleven chunks arrive and the last has one item.
2. Measure throughput against chunk size¶
The right chunk size is where the per-call overhead stops mattering. Measure it on your sink:
import asyncpg
import time
async def rows_per_second(pool: asyncpg.Pool, rows: list[tuple], size: int) -> float:
async with pool.acquire() as conn:
t = time.perf_counter()
for batch in chunks(rows, size):
await conn.executemany("insert into items values ($1, $2, $3)", batch)
return len(rows) / (time.perf_counter() - t)
Results on a local PostgreSQL 17 container with an unlogged three-column table:
| Method | Chunk size | Rows per second |
|---|---|---|
executemany |
1 | 8,648 |
executemany |
100 | 196,871 |
executemany |
1,000 | 306,279 |
executemany |
10,000 | 400,527 |
copy_records_to_table |
1,000 | 477,250 |
copy_records_to_table |
10,000 | 1,543,991 |
The first jump — 1 to 100 — is twentyfold; beyond 1,000 the gains are incremental. COPY is a different protocol path entirely, streaming rows without per-row statement overhead, and is the method to reach for when the sink supports it. A logged table, indexes and a real network will lower every absolute number; the shape of the curve is what transfers.
Verify: measure at least four chunk sizes on your sink and pick the knee, not the maximum.
3. Bound chunks by bytes as well as count¶
A chunk of 1,000 small rows is a few hundred kilobytes; a chunk of 1,000 documents with embedded blobs can be hundreds of megabytes, enough to exceed request limits or exhaust memory. Bound by both:
async def achunks_bounded(source, max_items: int, max_bytes: int, size_of=len):
batch: list = []
nbytes = 0
async for item in source:
n = size_of(item)
if batch and (len(batch) >= max_items or nbytes + n > max_bytes):
yield batch
batch, nbytes = [], 0
batch.append(item)
nbytes += n
if batch:
yield batch
Set max_bytes from the sink's limits — many bulk APIs cap request bodies at a few megabytes, Postgres has no hard limit on COPY but memory does — and from what you are willing to buffer per in-flight chunk. An item larger than max_bytes still goes through, alone in its own chunk; reject it explicitly if the sink would.
Verify: feed a mix of small and very large items; no chunk exceeds max_bytes except a single oversized item.
4. Run chunks concurrently, bounded by the sink¶
Once writes are chunked, concurrency across chunks adds throughput until the sink saturates. Bound it by the connection pool:
async def load_concurrently(pool: asyncpg.Pool, rows, size: int = 5000, parallel: int = 4) -> None:
sem = asyncio.Semaphore(parallel)
async def load(batch):
async with sem, pool.acquire() as conn:
await conn.copy_records_to_table("items", records=batch)
async with asyncio.TaskGroup() as tg:
for batch in chunks(rows, size):
await sem.acquire() # don't create all tasks up front
sem.release()
tg.create_task(load(batch))
Measured: COPY chunks of 5,000 over four connections reached 2,833,049 rows per second, against 1,543,991 for chunks of 10,000 on one connection. Waiting on the semaphore before creating each task matters for large inputs: creating one task per chunk up front would materialise the whole input as pending chunks in memory. Past the database's capacity — CPU, WAL bandwidth, lock contention on the table — more parallelism slows everything; measure with 2, 4, 8 and stop at the knee.
Verify: throughput rises with parallel up to a point and then flattens; pick the lowest value near the top.
5. Make chunk failures retryable¶
A chunk is also the unit of failure. If chunk 417 fails, you want to retry chunk 417, not restart the load. Two properties make that possible: each chunk is written in its own transaction, and the write is idempotent so a retry after an ambiguous failure — a timeout after the commit — does not duplicate rows:
async def load_chunk(pool, batch, attempts: int = 3) -> None:
for attempt in range(1, attempts + 1):
try:
async with pool.acquire() as conn, conn.transaction():
await conn.execute("create temp table stage (like items) on commit drop")
await conn.copy_records_to_table("stage", records=batch)
await conn.execute("insert into items select * from stage on conflict (id) do nothing")
return
except (asyncpg.PostgresConnectionError, OSError, TimeoutError):
if attempt == attempts:
raise
await asyncio.sleep(0.5 * 2 ** attempt)
Copying into a temporary staging table and inserting with ON CONFLICT DO NOTHING keeps COPY's speed while making the chunk idempotent. Record the last committed chunk as a checkpoint so a restarted job skips everything already loaded, as in checkpointing progress in long-running async jobs.
Verify: inject a failure on one chunk; it is retried, the load completes, and the row count equals the input size exactly.
Verification¶
Chunked processing is right when:
- The input is never fully materialised; chunking is lazy and the final partial chunk is kept.
- Chunk size sits at the measured knee for your sink, bounded by bytes.
- Concurrency across chunks is bounded by the sink's capacity, measured.
- Each chunk commits on its own and is idempotent, so failures retry one chunk.
Diagnostic Hook: export rows per second, chunk duration and chunk retries. A chunk duration that grows during a long load — common as indexes grow or autovacuum competes — means the sink is slowing; reduce parallelism before timeouts start. Retries concentrated on particular chunks point at bad data rather than infrastructure.
Pitfalls & edge cases¶
- Dropping the last partial chunk. Always flush the remainder.
- Chunking by count only. Large items produce oversized requests.
- Unbounded task creation per chunk. Materialises the input as pending tasks.
- One giant transaction for the whole load. A failure at 99% rolls everything back.
Frequently Asked Questions¶
What chunk size should I use for bulk inserts with asyncpg?
Measure, but expect most of the gain between 100 and 1,000 rows per call: in testing, executemany rose from 8,648 rows per second at one row to 306,279 at 1,000 and 400,527 at 10,000. COPY reached 1.5 million at 10,000.
How do I split an async iterator into chunks?
Append items to a list as they arrive and yield it when it reaches the chunk size, then yield the final partial list at the end. Add a byte limit, and a time limit if the stream can pause.
Should I insert chunks concurrently?
Yes, up to the database's capacity: COPY chunks over four connections reached 2.8 million rows per second against 1.5 million on one. Bound concurrency with a semaphore and stop increasing it at the throughput knee.
How do I retry a failed chunk without duplicating rows?
Write each chunk in its own transaction and make it idempotent — for example COPY into a temporary table and insert with ON CONFLICT DO NOTHING — then retry only that chunk.
Related¶
- Async Data Pipelines — up to the topic overview.
- Choosing chunksize for ProcessPoolExecutor work — the same trade-off for CPU work.
- Concurrent Execution & Worker Patterns — the section overview.