Bulk Loading Rows with asyncpg COPY¶
Loading data row by row is the slowest thing you can ask PostgreSQL to do, and from asyncio it is easy to do by accident: a loop of await conn.execute("insert ...") looks harmless. Measured with asyncpg 0.31 against PostgreSQL 17 on the same host, inserting 10,000 rows one statement at a time in autocommit took 8.8 s — each insert its own transaction, each commit its own disk flush. executemany loaded 1,000,000 rows in 3.9 s, and copy_records_to_table, which uses PostgreSQL's COPY protocol, in 1.7 s. This guide uses COPY for bulk loads, streams records into it, upserts through a staging table, and handles the all-or-nothing failure behaviour.
Prerequisites¶
- Python 3.11+,
pip install asyncpg; measured with asyncpg 0.31 and PostgreSQL 17. - asyncpg pools and transactions, from running transactions safely with asyncpg pools.
- Event-loop blocking, from avoiding event loop blocking with asyncpg.
1. Stop inserting one row per statement¶
The cost of a single-row insert in autocommit is dominated by fixed overhead: a round trip, a transaction, and a write-ahead-log flush on commit. Batch the rows:
# slow: one transaction and one round trip per row
for row in rows:
await conn.execute("insert into events values ($1, $2, $3, $4)", *row) # 10k rows: 8.8 s
# better: asyncpg pipelines executemany in one transaction
await conn.executemany("insert into events values ($1, $2, $3, $4)", rows) # 1M rows: 3.9 s
executemany in asyncpg sends the statements without waiting for each result and runs them in one implicit transaction, which removes both the round trips and the per-row commits. At 1 million rows it took 3.9 s; the row-by-row loop, at 0.88 ms per row, would have taken roughly 15 minutes. If you cannot batch — rows arrive one at a time — at least wrap a group of inserts in async with conn.transaction() so they share one commit.
Verify: time your current load path on 10,000 rows; anything near a millisecond per row is paying per-row commits.
2. Use copy_records_to_table for bulk loads¶
COPY streams rows in PostgreSQL's binary format with no per-statement parsing or planning. asyncpg exposes it for Python tuples:
async def load(pool, rows: list[tuple]) -> None:
async with pool.acquire() as conn:
await conn.copy_records_to_table(
"events",
records=rows,
columns=["id", "name", "amount", "ts"], # match the tuple order
)
Measured: 1,000,000 rows in 1.71 s, against 3.91 s for executemany. Values are encoded by asyncpg's type codecs, so they must already be the right Python types — Decimal for numeric, aware datetime for timestamptz; a wrong type raises in the client before anything is sent (a string in a numeric column raised decimal.InvalidOperation). Name columns explicitly so a later schema change does not silently shift values into the wrong columns.
Verify: load a sample with copy_records_to_table and compare row counts and a checksum with the source.
3. Stream records instead of building a list¶
A million tuples in a list is hundreds of megabytes. records also accepts an async iterable, so you can stream from a file, an API or another query, holding only a small window in memory:
async def rows_from_csv(path: str):
async with aiofiles.open(path) as f:
async for line in f:
id_, name, amount, ts = line.rstrip("\n").split(",")
yield int(id_), name, Decimal(amount), datetime.fromisoformat(ts)
async def load_csv(pool, path: str) -> None:
async with pool.acquire() as conn:
await conn.copy_records_to_table("events", records=rows_from_csv(path))
Measured on 100,000 rows, the async generator was as fast as the list (0.19 s against 0.21 s): the bottleneck is encoding and the server, not iteration. For data that is already in CSV form, copy_to_table with source= a file path or file-like object skips Python-side parsing entirely and lets PostgreSQL parse the text.
Verify: memory stays flat while loading a file many times larger than available RAM.
4. Upsert through a staging table¶
COPY only inserts; it has no ON CONFLICT. To merge a load into a table with existing rows, copy into a temporary staging table and merge with one INSERT ... SELECT:
async def upsert(conn, rows) -> str:
async with conn.transaction():
await conn.execute("create temp table stage (like events including defaults) on commit drop")
await conn.copy_records_to_table("stage", records=rows)
return await conn.execute("""
insert into events select * from stage
on conflict (id) do update
set name = excluded.name, amount = excluded.amount, ts = excluded.ts
""")
Measured: upserting 100,000 rows that all conflicted with existing keys took 1.61 s, returning INSERT 0 100000. The temporary table lives only inside the transaction (on commit drop) and is private to the connection, so concurrent loads do not collide. The same pattern does deletes, deduplication or validation with SQL before the merge.
Verify: run the upsert twice with the same data; the table's row count does not change and the values match the latest load.
5. Plan for all-or-nothing failure¶
A COPY is one statement. If any row fails a constraint, the whole statement fails and nothing is loaded:
try:
await conn.copy_records_to_table("events", records=rows)
except asyncpg.UniqueViolationError as exc:
log.error("load rejected: %s", exc.detail) # e.g. Key (id)=(10) already exists.
# zero rows were written
Measured: a load of 60,001 rows with one duplicate key at position 50,000 raised UniqueViolationError and left the table with 0 rows. That is usually what you want for a batch — no half-loaded state. For large loads, split them into chunks of, say, 100,000 rows, each its own COPY, and record which chunks succeeded so a retry resumes rather than repeats, as in checkpointing progress in long-running async jobs. When individual bad rows must be skipped, load into a staging table without constraints and filter with SQL.
Verify: inject a bad row into a test load; the error names it, and the target table is unchanged.
Verification¶
Bulk loading is set up well when:
- No load path commits per row; batches share a transaction.
- Large loads use COPY, with explicit columns and correctly typed values.
- Upserts go through a staging table in one transaction.
- Failures leave no partial state, and large loads are chunked and resumable.
Diagnostic Hook: log rows per second for each load job. Rates around a thousand rows per second mean per-row commits; hundreds of thousands mean COPY is working. A rate that falls as the table grows points at indexes and constraints on the target — load into a staging table, or drop and rebuild secondary indexes for very large initial loads.
Pitfalls & edge cases¶
- Row-by-row inserts in autocommit. Measured at 0.88 ms per row, dominated by commits.
- Untyped values. COPY encodes by column type; convert to
Decimal,datetimeand so on first. - Expecting ON CONFLICT from COPY. Use a staging table.
- One huge COPY. A single bad row fails the whole thing; chunk large loads.
Frequently Asked Questions¶
What is the fastest way to insert many rows with asyncpg?
copy_records_to_table, which uses PostgreSQL's COPY. In testing it loaded 1 million rows in 1.7 s, against 3.9 s for executemany and 8.8 s for just 10,000 rows inserted one at a time.
Can asyncpg COPY do an upsert?
Not directly. Copy into a temporary staging table and run INSERT ... SELECT ... ON CONFLICT DO UPDATE from it in the same transaction.
What happens if one row fails during COPY?
The whole COPY fails and no rows are written. In testing, one duplicate key in 60,001 rows left the table empty.
Can I stream rows into asyncpg COPY without a list?
Yes. records accepts an async iterable, so an async generator can feed it from a file or another source with flat memory.
Related¶
- Async Database Drivers — up to the topic overview.
- Streaming large result sets with asyncpg cursors — the reverse direction: reading large tables without buffering.
- Network I/O & Protocol Handling — the section overview.