Skip to content

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

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.

Rows per second with input as fast as possible 3 horizontal bars comparing per-row INSERT x8 with the others. Rows per second with input as fast as possible per-row INSERT x8 4,614/s executemany, 500 88.9-91.1k/s COPY, 500 98-152k/s asyncpg 0.31.0, Postgres 17 in Docker on the same host; 100-byte rows; two runs. Batching moved the ceiling by 20-30x.

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.

Latency from produced to committed, by input rate A grid of 3 rows by 4 columns. Latency from produced to committed, by input rate sink max rate: p50 2,000/s: p50 / p99 100/s: p50 / p99 per-row INSERT x8 2,159 ms (saturated) 1.8 / 3.0 ms 1.5 / 3.2 ms COPY, 500 or 50 ms 64-101 ms 31 / 58 ms 33 / 58 ms COPY, 500, no time limit 64-76 ms 135 / 260 ms 2,512 / 4,963 ms The time limit bounds latency when input is slow; the size limit, when it is fast.

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.

One bad row in a batch of 500 A grid of 3 rows by 5 columns. One bad row in a batch of 500 strategy rows written dead-lettered round trips time write the batch, give up 0 - 1 1.9 ms retry row by row 499 1 500 623.0 ms bisect 499 1 19 22.4 ms Bisection costs about 2 x log2(500) writes per bad row.

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.

How should this sink write? A decision on What does the sink need with 4 outcomes. How should this sink write? What does the sink need? input always below per-row ceiling per-row INSERT 1.5-1.8 ms latency bursts or backfills batch: 500 rows or 50 ms 98-152k rows/s upserts executemany or staging table + merge COPY has no ON CONFLICT bad rows possible bisect on data errors 19 trips, not 500 And always: flush the open batch on the way out.

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.