Batching Database Writes in Async Pipelines¶
The last stage of many pipelines writes each item to a database, and a round trip per row caps the whole pipeline at whatever the database can do one statement at a time. Batching removes the cap — but a batch that waits to fill trades throughput for latency, and a batch with one bad row fails as a unit. Measured on Python 3.14 with asyncpg 0.31.0 against Postgres 17 in a local container, writing 100-byte rows: eight workers doing one INSERT each managed 4,614 rows/s. A single sink that collected rows into batches of up to 500 wrote 88,869–91,107 rows/s with executemany and 98,047–151,572 rows/s with COPY across two runs. The time limit is what keeps latency bounded when input slows: at 100 rows/s, a batch that flushed only when full held rows for a median of 2,512 ms and up to 4,962 ms; flushing after at most 50 ms kept the median at 33 ms and p99 at 57 ms. When one row in a 500-row batch violated a constraint, the whole COPY failed; retrying row by row took 500 round trips and 623 ms, bisecting took 19 and 22.4 ms. This guide builds that sink.
Prerequisites¶
- Python 3.11+ and asyncpg with a connection pool.
- Bulk loading basics, from chunking large inputs for async batch processing.
- The topic overview, Async Data Pipelines.
1. Measure what per-row writes cost¶
Start from the baseline: a pool of workers, each taking an item from the queue and inserting it:
async def per_row_writer(q, pool):
while (row := await q.get()) is not END:
async with pool.acquire() as conn:
await conn.execute("INSERT INTO events VALUES ($1, $2, $3)", *row)
Measured with eight such workers and input as fast as the producer could go: 4,614 rows/s, and because the input queue filled up, rows waited about 2.2 s between being produced and being committed. At lower input rates the same writers kept up easily — 1.8 ms median latency at 2,000 rows/s and 1.5 ms at 100 rows/s. That is the trade-off batching will change: per-row writes have the lowest latency when the database keeps up, and the lowest ceiling when it does not. A pipeline sink that has to absorb bursts or backfills needs the higher ceiling.
Verify: you know your per-row ceiling and the input rate your pipeline must sustain.
2. Flush on size or time, whichever comes first¶
A batching sink collects items until the batch is full or the oldest item has waited long enough, then writes. The time limit must start when the first item of a batch arrives, not when the last one did:
END = object()
async def batch_writer(q, write, max_size=500, max_wait=0.05):
loop = asyncio.get_running_loop()
batch, deadline = [], None
while True:
timeout = None if not batch else max(0.0, deadline - loop.time())
try:
async with asyncio.timeout(timeout):
item = await q.get()
except TimeoutError:
item = None # the oldest item has waited long enough
if item is END:
if batch:
await write(batch)
return
if item is not None:
if not batch:
deadline = loop.time() + max_wait # clock starts with the first item
batch.append(item)
if batch and (len(batch) >= max_size or item is None):
await write(batch)
batch = []
With an empty batch, the writer waits for input with no timeout, so an idle pipeline costs nothing. Measured at three input rates with batches of 500: at full speed, flushes were driven by size and latency was 64–109 ms; at 2,000 rows/s, the 50 ms limit fired before batches filled and latency was 31 ms median, 58 ms p99; at 100 rows/s, 33 ms and 57 ms. Without a time limit, the same sink at 2,000 rows/s held rows 135 ms median and 259 ms p99 — the time to collect 500 rows — and at 100 rows/s, 2.5 s and 5.0 s.
Verify: at your lowest realistic input rate, p99 write latency stays below max_wait plus one write's duration.
3. Prefer COPY for plain inserts¶
Once rows are batched, how the batch is sent matters. executemany sends one parameterised statement per row in a single pipelined exchange; COPY streams all rows in Postgres's bulk-load format:
async def write_copy(batch):
async with pool.acquire() as conn:
await conn.copy_records_to_table(
"events", records=batch, columns=["id", "payload", "created"]
)
Measured: COPY wrote 98,047–151,572 rows/s against 88,869–91,107 for executemany with the same batches; the spread between runs reflects other load on the shared machine. COPY cannot express ON CONFLICT or RETURNING; for upserts, copy into a temporary table and merge with one INSERT ... SELECT ... ON CONFLICT per batch, or use executemany with the upsert statement. Size batches so a single write stays short — here 500 rows took a few milliseconds — because a write holds a pool connection and blocks the sink for its whole duration, and very large batches turn every flush into a latency spike.
Verify: throughput with your batch size and write method is measured, and a single flush takes a small fraction of max_wait.
4. Isolate bad rows without retrying the whole batch row by row¶
A constraint violation in one row fails the entire COPY or executemany, and the batch's other rows are not written. Dropping the batch loses good data; retrying each row individually throws away the batching. Split the batch in half and retry each half, recursively, until the failures are single rows:
async def write_or_bisect(batch, dead: list) -> None:
try:
await write_copy(batch)
except (asyncpg.IntegrityConstraintViolationError, asyncpg.DataError) as exc:
if len(batch) == 1:
dead.append((batch[0], repr(exc)))
return
mid = len(batch) // 2
await write_or_bisect(batch[:mid], dead)
await write_or_bisect(batch[mid:], dead)
Measured with one row in 500 violating a CHECK constraint: the plain COPY raised CheckViolationError and wrote nothing; row-by-row retry wrote 499 rows and dead-lettered 1 in 500 round trips and 623 ms; bisection did the same in 19 round trips and 22.4 ms. The cost grows with the number of bad rows — each one costs about 2 × log₂(batch size) extra writes — so a batch full of bad data should be dead-lettered whole after a limit rather than bisected. Catch only data errors here; a lost connection or a full disk is a pipeline error and should stop the sink, as described in handling per-item errors in async pipelines.
Verify: a batch containing one invalid row writes every valid row and dead-letters exactly the invalid one.
5. Flush the last batch on shutdown and cancellation¶
A batching sink holds data that the rest of the pipeline considers delivered. If the sink is cancelled — a deploy, a shutdown signal, a TaskGroup aborting — the open batch is lost unless the sink writes it on the way out:
async def batch_writer(q, write, max_size=500, max_wait=0.05):
batch = []
try:
... # the loop from step 2
finally:
if batch:
await asyncio.shield(write(batch)) # survives a second cancel
Measured: with 1,300 rows queued and the sink cancelled 0.2 s later with a 10 s time limit, the version without the finally wrote 1,000 rows — two full batches — and lost the 300 in the open batch; with it, all 1,300 were written. The shield keeps a second cancellation from interrupting the final write; bound the shutdown overall with a deadline, as in enforcing a hard shutdown deadline. If upstream acknowledgements (a message broker's acks, a checkpoint) are tied to rows, send them only after the batch containing those rows has committed.
Verify: cancelling the pipeline with a partly filled batch writes every row that reached the sink.
Verification¶
A batching sink is correct when:
- It flushes on size or time, with the time limit measured from the first item of the batch.
- Throughput and latency are measured at the fastest and slowest input rates.
- Data errors are bisected to the offending rows, and other errors stop the pipeline.
- The open batch is written on cancellation and shutdown, before any upstream acknowledgement.
Diagnostic Hook: record each flush's size and the reason — size or time. Mostly size-triggered flushes mean the sink is running near its ceiling and latency is set by fill time; mostly time-triggered flushes of small batches mean the input is slow and max_wait is what users feel. A sudden mix change is the earliest sign that input volume has shifted.
Pitfalls & edge cases¶
- Size-only flushing. Measured: 5.0 s p99 latency at 100 rows/s.
- Retrying failed batches row by row. Measured: 500 round trips against 19 for bisection.
- Losing the open batch on cancel. Measured: 300 of 1,300 rows gone.
- Batches so large that one flush is a latency spike. Keep a single write short.
Frequently Asked Questions¶
How do I batch database inserts in an asyncio pipeline?
Have one sink task collect items from the queue and write when the batch reaches a size limit or its oldest item reaches a time limit. With 500 rows or 50 ms, COPY wrote 98,000 to 152,000 rows/s in testing, against 4,614 for per-row inserts.
What batch size and flush interval should I use?
Choose the time limit from the latency you can accept and the size so that one write stays short. At 100 rows/s, a size-only batch held rows for up to 5 s; a 50 ms limit kept p99 at 57 ms.
Is COPY faster than executemany in asyncpg?
For plain inserts, yes: 98,000 to 152,000 rows/s against about 90,000 with the same batches in testing. COPY cannot do ON CONFLICT, so use executemany or a staging table for upserts.
What happens when one row in a batch fails?
The whole batch fails. Bisect it: one bad row in 500 took 19 round trips to isolate, against 500 for row-by-row retries.
Related¶
- Async Data Pipelines — up to the topic overview.
- Scaling a slow pipeline stage — when the sink is the bottleneck.
- Concurrent Execution & Worker Patterns — the section overview.