Handbooks / Python asyncio / Chapter 5

Async Databases

42 pages · ~84 min✓ Reviewed

Builds on Async HTTP Clients & Servers. Next up: Async WebSockets & Streaming.

Part 1 · Async Databases in Python: Postgres, SQLite and Redis Without Blocking the Loop

Async Databases in Python: Postgres, SQLite and Redis Without Blocking the Loop

Almost every web service spends most of its time waiting on a database. In an asyncio program that wait is dangerous in one specific way: the event loop runs on a single thread, so one blocking call from psycopg2, sqlite3 or the synchronous redis-py client freezes every other coroutine until it returns. The tests pass, the demo feels fast, and then the server stalls under real load. An async driver avoids this because it yields at network I/O, letting the loop serve other tasks while Postgres or Redis does its work.

Going async does not make a single query faster. A 200 ms query is still 200 ms; the gain is that many requests can wait on I/O together. That also moves the bottleneck: the database and its connection pool, not the loop, now cap your throughput. Much of this chapter is about that second half, covering how pools are sized, why 1000 tasks do not mean 1000 parallel queries, and what goes wrong when sessions, transactions or timeouts are handled carelessly.

By the end you will be able to talk to Postgres with asyncpg (connections, pools, prepared statements) and with the SQLAlchemy 2.0 AsyncSession, choose a transaction isolation level and retry safely, and spot N+1 and lazy-loading traps before they reach production. You will also size a pool against max_connections, use aiosqlite and redis.asyncio correctly, and put deadlines on every call so that a timeout or cancellation leaves connections and transactions clean.

Before you start

You should be comfortable with async def, await, tasks and asyncio.gather, and know basic SQL. To follow along, use Python 3.11 or newer and install asyncpg, sqlalchemy[asyncio], aiosqlite and redis (5.0.1 or later for aclose()). A local Postgres and Redis are handy, for example from Docker, but SQLite needs nothing extra. While you experiment, run with asyncio.run(main(), debug=True) so the loop warns you about any callback that blocks for over 100 ms.

Part 2 · Why Async Database Access

One thread, one loop

An asyncio program runs every coroutine on a single thread, and the event loop decides who runs next. The loop can only switch tasks when the running task reaches an await. If a coroutine calls a blocking function such as cursor.fetchall() from psycopg2, sqlite3, or the synchronous redis-py client, the thread is stuck inside that call. Every other coroutine waits until it returns, including ones that have nothing to do with the database.

An async driver behaves differently. It sends the query over the socket and then awaits the reply. While it waits, the task is suspended and the loop is free to run something else. Only await hands control back to the loop, so a call that never awaits can never share the thread.

The program below makes the difference visible without needing a database. A background ticker adds one tick every 0.1 s. The blocking query uses time.sleep, which holds the thread like a sync driver does. The async query uses asyncio.sleep, which suspends like await conn.fetch(...) does. Both pretend to take 0.35 s.

python
import asyncio
import time

async def ticker(ticks):
    while True:
        await asyncio.sleep(0.1)
        ticks.append(1)

async def blocking_query():
    time.sleep(0.35)          # like cur.fetchall() in psycopg2

async def async_query():
    await asyncio.sleep(0.35)  # like await conn.fetch(sql)

async def ticks_during(query):
    ticks = []
    task = asyncio.create_task(ticker(ticks))
    await asyncio.sleep(0)     # let the ticker start
    await query()
    task.cancel()
    return len(ticks)

async def main():
    print('blocking query:', await ticks_during(blocking_query), 'ticks while it ran')
    print('async query:', await ticks_during(async_query), 'ticks while it ran')

asyncio.run(main())
output
blocking query: 0 ticks while it ran
async query: 3 ticks while it ran

The ticker stands in for every other request your server is handling. With the blocking call it gets no turns at all for 0.35 s. With the awaited call it keeps ticking. The sequence below is what happens each time you await a real driver call.

What happens at await conn.fetch(...)
  1. 1Task A awaitssends the query, then suspends
  2. 2Loop runs otherstasks B, C and D get their turns
  3. 3Postgres worksplans, scans, replies
  4. 4Task A resumesthe rows are returned to the caller

In a real handler it looks like the fragment below. The only change from a sync driver is that the call is awaited, and that single keyword is what lets the loop continue.

python
async def get_user(conn, uid):
    # conn is an asyncpg connection; awaiting suspends this task only
    return await conn.fetchrow('SELECT * FROM users WHERE id = $1', uid)

Fragment: needs a live asyncpg connection, so it is shown without running.

Concurrency, not speed

It is easy to expect async to make queries faster. It does not. A query that takes 200 ms inside Postgres still takes 200 ms when you await it. What async changes is what the rest of your program does during those 200 ms. Instead of one request occupying the thread, many requests can all be waiting on I/O at once.

Driver1 query10 queries at once
Sync200 msabout 2000 ms, one after another
Async200 msabout 200 ms, waiting together

The next program scales the numbers down to 0.1 s per query so it finishes quickly. A fake query awaits asyncio.sleep, standing in for the network round trip. We time one query, then ten awaited in a row, then ten started together with asyncio.gather.

python
import asyncio
import time

async def query():
    await asyncio.sleep(0.1)   # stand-in for a 100 ms query

async def timed(label, coro):
    start = time.perf_counter()
    await coro
    print(label, round(time.perf_counter() - start, 1), 's')

async def one_after_another():
    for _ in range(10):
        await query()

async def together():
    await asyncio.gather(*(query() for _ in range(10)))

async def main():
    await timed('one query:', query())
    await timed('ten, one at a time:', one_after_another())
    await timed('ten together:', together())

asyncio.run(main())
output
one query: 0.1 s
ten, one at a time: 1.0 s
ten together: 0.1 s

A single query is no faster than before. The gain only appears when many requests overlap. If one endpoint is slow because its query is slow, the fix is the query, an index, or a cache. Changing the driver will not help.

The database is still the limit

Starting a thousand coroutines does not give you a thousand parallel queries. Postgres serves them over a limited number of connections, and an asyncpg or SQLAlchemy pool hands those out. If 1000 coroutines share 10 connections, 10 run and 990 wait on the pool. A pool behaves like a semaphore, so we can model it with one. Here 20 tasks share a pool of 5.

python
import asyncio
import time

pool = asyncio.Semaphore(5)    # a pool of 5 connections
in_use = 0
peak = 0

async def request():
    global in_use, peak
    async with pool:           # acquire a connection, or wait
        in_use += 1
        peak = max(peak, in_use)
        await asyncio.sleep(0.1)   # the query
        in_use -= 1

async def main():
    start = time.perf_counter()
    await asyncio.gather(*(request() for _ in range(20)))
    print('peak connections in use:', peak)
    print('elapsed:', round(time.perf_counter() - start, 1), 's')

asyncio.run(main())
output
peak connections in use: 5
elapsed: 0.4 s

Twenty tasks at five per round take four rounds of 0.1 s. Waiting on the pool costs almost nothing, because a waiting coroutine uses no thread, so a small pool is normal. Growing it only helps while the database has spare capacity. Pool sizing is covered in its own section later in this chapter.

Concurrency is the win

Async lets many requests wait on I/O together. It does not shorten any one query, and it does not lift the limit set by the database and its connection pool.

Async drivers and the to_thread escape hatch

Each common database has a blocking driver and an async counterpart. Pick the async one for anything that runs inside async def.

DatabaseBlocking driverAsync driver
Postgrespsycopg2asyncpg (binary protocol) or psycopg 3 async
SQLitesqlite3aiosqlite (runs sqlite3 on a worker thread)
Redisredis-py sync clientredis.asyncio
MySQLPyMySQLaiomysql or asyncmy

The driver families differ in how they stay async. asyncpg and redis.asyncio speak their wire protocols over non-blocking sockets, so the loop does the waiting. aiosqlite is the odd one out. SQLite is an in-process library with no network, so aiosqlite runs the normal sqlite3 calls on a dedicated thread and awaits the result. The loop stays free in either case.

When only a sync library exists

Sometimes the library you need has no async version, such as a legacy ORM or a vendor SDK. In that case await asyncio.to_thread(fn, *args) runs the call in a worker thread and suspends your coroutine until it finishes. The loop thread never blocks. The program below runs three 0.2 s blocking calls at the same time.

python
import asyncio
import time

def legacy_lookup(uid):
    time.sleep(0.2)            # a sync-only library blocking
    return f'user-{uid}'

async def main():
    start = time.perf_counter()
    users = await asyncio.gather(
        *(asyncio.to_thread(legacy_lookup, i) for i in (1, 2, 3))
    )
    print(users)
    print('elapsed:', round(time.perf_counter() - start, 1), 's')

asyncio.run(main())
output
['user-1', 'user-2', 'user-3']
elapsed: 0.2 s

The three calls overlapped because each ran on its own thread, so the loop was free the whole time. This is an escape hatch, not a goal. The default thread pool is small, so a hot path that sends thousands of calls through to_thread will queue behind it. Use an async driver whenever one exists.

Which call goes in async def?

The same idea applies to other blocking calls. Here is what each one should become.

Blocking callAsync replacement
time.sleep()asyncio.sleep()
requests.get()httpx.AsyncClient
psycopg2asyncpg
reading a big file with open()asyncio.to_thread()

Mistakes and finding stalls

Common mistake: a sync driver or ORM inside async def

The code looks fine and unit tests pass, because a single request has nothing else to starve. Under load, each blocking call freezes the whole server for its duration. One slow sqlite3 call delays every other request, including health checks.

Common mistake: expecting async to speed up a slow query

The query takes the same time as before. Fix it with an index, a rewritten query or a cache. Async only helps with how many requests can wait at once.

Common mistake: treating 1000 tasks as 1000 parallel queries

They share a pool. Past the pool size the extra tasks queue, and raising the pool size does not help once Postgres is busy.

Blocking calls are hard to spot by reading code, so asyncio can report them for you. Run your program with asyncio.run(main(), debug=True), or set the environment variable PYTHONASYNCIODEBUG=1. In debug mode the loop logs a warning such as Executing <Task ...> took 0.212 seconds whenever one step of a task hogs the thread for more than 100 ms. The logged task shows where to look.

To catch shorter stalls, lower the threshold from inside the running loop with asyncio.get_running_loop().slow_callback_duration = 0.05. This only has an effect in debug mode, so use it in development and not in production.

The symptoms you see in production also point to the cause. Use this table to tell a blocked loop from an exhausted pool or a plain slow query.

SymptomLikely cause
All requests get slow at the same momentsA blocking sync call on the loop
p99 latency spikes under loadThe pool is exhausted and tasks are queuing
Only one endpoint is slowThe query itself
Rules for the rest of the chapter

Await every database call inside async def. For sync-only libraries use asyncio.to_thread. Run with debug=True during development, and load-test before release, because unit tests will not reveal loop stalls.

Part 3 · asyncpg Connections and Pools

One Connection and the Query API

asyncpg talks to Postgres over its own binary protocol and never blocks the event loop. The simplest way in is a single connection: await asyncpg.connect(dsn) opens one network connection to the server and gives you a connection object. That object is backed by a real Postgres backend process, so you must give it back with await conn.close(). Put the close in a finally block so it runs even when a query raises.

python
# uses the asyncpg package; this block only defines a function
async def count_users(dsn):
    conn = await asyncpg.connect(dsn)
    try:
        return await conn.fetchval('SELECT count(*) FROM users')
    finally:
        await conn.close()

An unclosed connection keeps its Postgres backend alive until the process exits.

Once you have a connection, four methods cover almost every query. They differ only in the shape of what they hand back, so pick the one that matches the result you expect.

MethodReturnsTypical use
fetch()A list of Record objects (empty list if no rows)Reading many rows
fetchrow()One Record, or None when nothing matchesLookup by id
fetchval()A single value from the first column of the first rowcount(*), a generated id
execute()A status string such as 'INSERT 0 1'Writes where you need no rows back

Placeholders keep values out of the SQL

asyncpg uses Postgres-style positional placeholders: $1, $2 and so on. The %s of psycopg2 and the ? of SQLite are not understood and give a syntax error. The values are passed as separate arguments after the SQL string and travel to the server apart from the SQL text, so the server never parses them as SQL. That is what blocks injection: a value like '; DROP TABLE users; -- is just data.

python
async def find_user(conn, uid, org):
    r = await conn.fetchrow(
        'SELECT * FROM users WHERE id = $1 AND org = $2',
        uid, org)
    if r is None:
        return None
    return r['name'], r[0], dict(r)

The same Record read by column name, by position and as a plain dict.

Records and native types

A Record behaves like a light mix of tuple and dict. You can read a column by name with r['name'], by position with r[0], or turn the whole row into a dictionary with dict(r). Values arrive as real Python objects, because asyncpg decodes the binary protocol directly rather than handing you strings.

Postgres typePython type you receive
uuiduuid.UUID
timestamptzdatetime.datetime
numericdecimal.Decimal
int[]list of int
Common mistakes

Writing %s or ? instead of $1 fails with a syntax error. Skipping close() leaks a server connection each time the code runs.

Pools: Sharing Connections Safely

Opening a connection costs a network round trip plus authentication, and Postgres only allows a limited number of them. A pool opens connections once and lends them out. You create it a single time with await asyncpg.create_pool(dsn, min_size=10, max_size=10); both sizes default to 10. Each unit of work then borrows a connection with async with pool.acquire() as conn: and the connection returns to the pool when the block ends, even if it raised.

python
async def top_orders(pool, uid):
    async with pool.acquire() as conn:
        return await conn.fetch(
            'SELECT * FROM orders WHERE user_id = $1 LIMIT 10', uid)

Inside the async with block you own that connection exclusively, so run every query that belongs together there. When you only need one query, the shortcuts pool.fetch(), pool.fetchrow(), pool.fetchval() and pool.execute() do the whole cycle in one call: acquire, run the query, release. Reach for them first, and use acquire() only when several queries must share one connection or one transaction.

What the pool does on release

A borrowed connection may have been changed by the previous user, so the pool cleans it before handing it to anyone else. It runs RESET ALL to undo session settings, UNLISTEN to drop notification subscriptions, and closes any cursors left open. Separately, max_inactive_connection_lifetime (300 seconds by default) closes connections that sit idle that long, and the pool opens new ones when demand returns.

Per-connection setup with init

Some setup belongs to each connection rather than to a query, such as teaching asyncpg how to decode jsonb. Pass an init= coroutine to create_pool and it runs once on every new connection before the pool uses it. By default asyncpg returns jsonb as a string; the codec below makes it return Python objects instead.

python
import json

async def init(conn):
    await conn.set_type_codec(
        'jsonb',
        encoder=json.dumps,
        decoder=json.loads,
        schema='pg_catalog')

async def start(dsn):
    return await asyncpg.create_pool(dsn, init=init)

A codec set inside a query would be erased by RESET ALL on release, so install it through init.

One pool per process

Create the pool once at startup and close it at shutdown. Creating a pool per request throws away the whole benefit.

One Query at a Time per Connection

A Postgres connection speaks a strict request and reply conversation, so a connection serves one query at a time, and every concurrent task needs its own connection. If two coroutines await queries on the same connection at the same moment, asyncpg raises InterfaceError: cannot perform operation: another operation is in progress. The program below uses stand-in classes instead of asyncpg, so it runs anywhere. FakeConn raises the same error when it is busy, and FakePool hands each task a free connection from a queue, which is what a real pool does.

python
import asyncio

class InterfaceError(Exception):
    pass

class FakeConn:
    def __init__(self):
        self.busy = False

    async def fetch(self, sql):
        if self.busy:
            raise InterfaceError('another operation is in progress')
        self.busy = True
        try:
            await asyncio.sleep(0.01)
            return [sql]
        finally:
            self.busy = False

class FakePool:
    def __init__(self, size):
        self.free = asyncio.Queue()
        for _ in range(size):
            self.free.put_nowait(FakeConn())

    async def fetch(self, sql):
        conn = await self.free.get()
        try:
            return await conn.fetch(sql)
        finally:
            self.free.put_nowait(conn)

async def shared():
    conn = FakeConn()
    try:
        await asyncio.gather(conn.fetch('q1'), conn.fetch('q2'))
    except InterfaceError as e:
        print('shared conn:', e)

async def pooled():
    pool = FakePool(2)
    rows = await asyncio.gather(pool.fetch('q1'), pool.fetch('q2'))
    print('pooled:', rows)

async def main():
    await shared()
    await pooled()

asyncio.run(main())
output
shared conn: another operation is in progress
pooled: [['q1'], ['q2']]

In shared(), the first fetch marks the connection busy and suspends at its sleep. The second fetch then starts on the same object and hits the error. In pooled(), each call takes a different connection from the queue, so both run side by side and both succeed. The fix in real code is the same: give each task its own connection through pool.fetch() or its own pool.acquire(), and never pass one connection to several tasks.

Why does gathering two fetches on one connection fail, while two pool.fetch() calls work?
  • Both fetch calls on a single connection start before either one finishes, and the connection can only hold one operation at a time, so the second call is rejected.
  • Each pool.fetch() call borrows its own connection, so the two queries never touch the same one and can wait on the network together.
  • Sharing a connection is only safe when the awaits run one after another, for example two await lines in a row.
# fails: two awaits at once on one conn
await asyncio.gather(conn.fetch(q1), conn.fetch(q2))

# works: each call gets its own connection
await asyncio.gather(pool.fetch(q1), pool.fetch(q2))
Common mistake

Passing one connection into several tasks, for example through TaskGroup or gather, raises InterfaceError as soon as two of them overlap.

Bulk Loads: executemany and COPY

Inserting rows one at a time in a loop costs a network round trip per row, and the waiting dominates the runtime. asyncpg offers two ways to cut that. executemany() sends one statement with many sets of arguments in a batch. copy_records_to_table() uses the Postgres COPY protocol, which streams all the rows in a single operation and is often 10x or more faster than row-by-row INSERT.

MethodRound tripsSpeedChoose it when
execute() in a loopOne per rowSlowestA handful of rows
executemany()BatchedFasterFew rows, or you need ON CONFLICT
copy_records_to_table()One COPY streamOften 10x or moreThousands of rows into one table
python
async def load_events(conn, rows):
    # rows is a list of tuples, e.g. [(1, 'login'), (2, 'logout')]
    async with conn.transaction():
        await conn.copy_records_to_table(
            'events', records=rows, columns=['user_id', 'kind'])

async def load_few(conn, rows):
    await conn.executemany(
        'INSERT INTO events(user_id, kind) VALUES($1, $2)', rows)

Pass columns=[...] to COPY into a subset of the table's columns.

Wrapping a load in conn.transaction() makes it all or nothing: if one row is bad, none of them are kept. COPY cannot express ON CONFLICT, so when rows might already exist, use executemany() with an upsert statement, or COPY into a staging table first and merge from there.

Key points

Close connections with finally or return them with async with. Run one query at a time per connection. Pass values as $1 parameters. Put per-connection setup in init, and load large batches with COPY.

Common mistake

Looping over thousands of rows with await conn.execute(INSERT ...) works, but it spends nearly all its time waiting on round trips. Switch to COPY.

Part 4 · Prepared Statements

Parse once, execute many times

Every SQL string you send to Postgres has to be parsed, analysed and planned before any row is read. A prepared statement splits that work from the execution. The server parses and plans the query once, keeps the result under a name, and each later call only sends new values for the $1, $2 placeholders. For a query you run thousands of times a second, skipping the repeated parse and plan work is a real saving.

Life of a prepared statement
  1. 1ParseSQL text checked once
  2. 2Planserver picks a strategy
  3. 3Execute with $1first set of values
  4. 4Execute with $1next values, no reparse

You rarely have to ask for this in asyncpg, because the driver prepares every query automatically. When you call conn.fetch(), fetchrow() or execute() with a SQL string, asyncpg prepares it, runs it, and remembers the prepared statement in a cache that belongs to that one connection. The next time the same SQL text goes through the same connection, asyncpg skips the parse step and just executes.

The cache holds statement_cache_size=100 statements by default and behaves as an LRU: when a 101st distinct query arrives, the statement that was used least recently is thrown away. Because the cache lives on the connection, a pool does not share it. Ten pooled connections means ten separate caches, and each one warms up on its own. The small simulation below shows the eviction rule with a cache of two entries.

python
from collections import OrderedDict

class StatementCache:
    def __init__(self, size):
        self.size = size
        self.items = OrderedDict()
        self.prepares = 0

    def get(self, sql):
        if sql in self.items:
            self.items.move_to_end(sql)
            return 'hit'
        self.prepares += 1
        self.items[sql] = f'stmt_{self.prepares}'
        if len(self.items) > self.size:
            self.items.popitem(last=False)
        return 'prepare'

cache = StatementCache(size=2)
for sql in ['A', 'B', 'A', 'C', 'B']:
    print(sql, cache.get(sql))
print('prepares:', cache.prepares)
output
A prepare
B prepare
A hit
C prepare
B prepare
prepares: 4

Query B had to be prepared a second time. Adding C pushed out B, which was the least recently used entry at that moment. In a real application the same thing happens if you build more than 100 distinct SQL strings per connection, for example by pasting values into the text instead of using placeholders.

When you want to hold on to a statement yourself, for instance to run it in a tight loop, prepare it explicitly. The returned object has the same fetch methods as a connection, minus the SQL argument.

python
stmt = await conn.prepare("SELECT * FROM users WHERE id=$1")

r1 = await stmt.fetchrow(42)
r2 = await stmt.fetchrow(7)

The second call sends only the value 7. Keep in mind that stmt is tied to the connection that created it. If you prepare on one connection from a pool and later run the object after releasing that connection, you are no longer in control of who else is using it, so prepare and use the statement inside a single async with pool.acquire() block.

Custom plans, generic plans and skewed data

Planning depends on the values in the query. Postgres can look at $1, compare it with table statistics, and build a plan for that exact value. That is a custom plan. A generic plan is built without looking at the value, so it can be reused for every call. Postgres tries custom plans for the first five executions of a prepared statement, then compares the generic plan's estimated cost with the average of those custom plans. If the generic plan looks about as cheap, it may switch to it from the sixth execution onward to save planning time.

Execution or settingPlan usedWhy
Runs 1 to 5CustomThe planner sees the real $1 value
Run 6 onward (auto)Maybe genericUsed only when its cost is close to the custom plans
plan_cache_mode = force_custom_planAlways customReplans every time, which suits skewed data
plan_cache_mode = force_generic_planAlways genericSkips replanning completely

The trouble starts when the data is lopsided. Imagine an orders table where status = 'rare' matches 50 rows and status = 'common' matches nine million. With the real value in view, the planner uses an index scan for 'rare' and a sequential scan for 'common'. A generic plan has no value to look at, so it commits to one shape for both. If that shape is the sequential scan, your rare lookups become slow; if it is the index scan, the common lookups crawl through millions of random reads.

status valueRows matchedBest planGeneric plan
'rare'50Index scanOne plan for both
'common'9 millionSequential scanWrong for one of them

The symptom is a query that runs fast for a while and then becomes slow after roughly the sixth call on a connection, only for some parameter values. The fix is to tell Postgres never to switch: set plan_cache_mode to force_custom_plan. With asyncpg you can send it as a server setting so that every pooled connection gets it when it connects.

python
pool = await asyncpg.create_pool(
    dsn,
    server_settings={"plan_cache_mode": "force_custom_plan"},
)

You pay a little planning time on each call in exchange for a plan that matches the value. For queries on evenly distributed columns the default auto is fine and you can leave it alone.

Common mistake

Trusting the generic plan on a skewed column. If one query is slow only after several runs, or only for certain values, compare EXPLAIN (ANALYZE) output for a fresh connection against a warmed one before blaming the index.

Prepared statements behind PgBouncer

Prepared statements live inside one Postgres server connection, which is why they clash with connection poolers. PgBouncer in transaction mode hands a server connection to a client only for the length of one transaction, then gives it to someone else. asyncpg, however, prepares a statement and later executes it by name, assuming the same server connection is still behind its own socket. After PgBouncer moves you, that name is unknown.

The message names a statement such as __asyncpg_stmt_1__, the internal name asyncpg gave it. It often appears only under load, since with a quiet pooler you may keep getting the same server connection and everything looks fine in testing.

The classic fix is to turn the client-side statement cache off, so asyncpg never relies on a name that must survive across transactions. Set statement_cache_size=0 when connecting.

python
conn = await asyncpg.connect(dsn, statement_cache_size=0)

pool = await asyncpg.create_pool(dsn, statement_cache_size=0)

Through SQLAlchemy there are two caches to switch off. One is asyncpg's own cache; the other is a cache the SQLAlchemy asyncpg dialect keeps of its prepared statement names. Pass both through connect_args, otherwise the error can still come back.

python
engine = create_async_engine(
    URL,
    connect_args={
        "statement_cache_size": 0,
        "prepared_statement_cache_size": 0,
    },
)

Turning the cache off has a cost: every call is parsed again. A newer option keeps the performance. PgBouncer 1.21 and later can track protocol-level prepared statements itself and replay them on whichever server connection a client lands on. Enable it by setting max_prepared_statements to a value above zero in pgbouncer.ini; once it is on you can leave asyncpg's statement cache at its default.

python
; pgbouncer.ini
pool_mode = transaction
max_prepared_statements = 100
SetupWhat to doCost
Direct to PostgresKeep the default cacheNone
PgBouncer before 1.21, transaction modestatement_cache_size=0 (plus prepared_statement_cache_size=0 in SQLAlchemy)Every query is parsed again
PgBouncer 1.21+, transaction modeSet max_prepared_statements above 0 and keep the cacheSmall memory use in PgBouncer
Common mistake

Leaving the statement cache on behind a transaction-mode PgBouncer that cannot track prepared statements. In SQLAlchemy, setting only statement_cache_size and forgetting prepared_statement_cache_size is the usual half-fix.

Schema changes while statements are cached

A cached prepared statement is a promise about the shape of the table at the time it was planned. If a deploy runs ALTER TABLE and changes a column type or adds a column that a SELECT * would now return, Postgres refuses to run the old plan and asyncpg raises InvalidCachedStatementError ("cached plan must not change result type").

What happens next depends on where the call was made. Outside a transaction, asyncpg catches the error, throws away the stale statement, prepares it again and retries once, so your code never notices. Inside a transaction that is not possible: statements earlier in the transaction may have already run, and replaying only the failing one would silently break atomicity. asyncpg therefore lets the exception reach you, and you must retry the whole transaction yourself.

python
from asyncpg.exceptions import InvalidCachedStatementError

async def fetch_user(pool, uid):
    async with pool.acquire() as conn:
        for attempt in range(2):
            try:
                async with conn.transaction():
                    return await conn.fetchrow(
                        "SELECT * FROM users WHERE id=$1", uid)
            except InvalidCachedStatementError:
                if attempt == 1:
                    raise

The failed attempt invalidates the cached statement, so the second pass prepares it afresh against the new table definition. One retry is enough; if it fails again something else is wrong and the error should surface. Naming the columns in the SELECT instead of using * also makes the problem much less likely, because adding a column then no longer changes the result type.

Common mistake

Assuming asyncpg handles every InvalidCachedStatementError for you. It does so only for statements outside a transaction. Inside conn.transaction() the error reaches your code, and a deploy that alters tables will show up as a short burst of failed requests unless you retry.

Why does the same code work in development and fail behind PgBouncer in production?
  • In development you connect straight to Postgres, so the server connection that holds the prepared statement is always the one asyncpg talks to.
  • In production a transaction-mode PgBouncer can give the next statement to another server connection that never saw the PREPARE, which produces the "does not exist" error.
  • Either disable the statement caches, or upgrade to PgBouncer 1.21+ and set max_prepared_statements above 0.
What to remember

asyncpg prepares and caches every query per connection, which is good for speed. Watch for three side effects: generic plans on skewed data (use force_custom_plan), transaction-mode poolers (turn the cache off or use PgBouncer 1.21+), and schema changes during deploys (retry inside transactions).

Part 5 · SQLAlchemy 2.0 Async Engine and AsyncSession

The engine and the session factory

SQLAlchemy 2.0 ships with a full asyncio layer, so you can keep the ORM you already know and still never block the event loop. Everything starts with an async engine. You build it with create_async_engine, and the URL names the async driver after the plus sign. postgresql+asyncpg tells SQLAlchemy to use the asyncpg dialect from the earlier sections, and pool_size=5 keeps five connections open in the engine's pool.

python
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine

engine = create_async_engine(
    "postgresql+asyncpg://user:pw@host/db",
    pool_size=5,
)

Session = async_sessionmaker(engine, expire_on_commit=False)

Both objects are created at module level or in startup code, not inside a request.

The engine is not a connection. It is the owner of the pool, and it only borrows a connection when a session actually needs to talk to the database. The second object, async_sessionmaker, is a session factory: calling Session() hands you a fresh AsyncSession bound to that engine. Build it once at startup and share it across the whole application, because it holds configuration and no per-request state.

The expire_on_commit=False option deserves attention. By default a commit marks every loaded object as expired, so the next read of user.name triggers a reload from the database. In a sync session that reload happens silently. In an async session it cannot be awaited from a plain attribute access, so it fails. Turning expiry off keeps the values you already loaded readable after the commit.

Three objects are involved, and they live for very different lengths of time. Mixing them up is the root of most async ORM bugs, so it helps to keep this table in mind.

ObjectCreatedShared?
AsyncEngineOnce per process, at startupYes, it owns the connection pool
async_sessionmakerOnce per process, at startupYes, it is only a factory
AsyncSessionOne per request or taskNever, it holds one unit of work

Running queries with 2.0 style

Most of the time you open a session and a transaction together. The idiom async with Session() as s, s.begin(): does both in one line. The async with on the session guarantees it is closed and its connection returns to the pool. The s.begin() part wraps your work in a transaction: it commits when the block ends cleanly and rolls back if an exception escapes.

What s.begin() does when the block ends

Inside the block you write queries in the 2.0 style: build a statement with select(), then pass it to execute(). The call returns a result object, and scalars() unwraps each row into the ORM object itself, so all() gives you a plain list of User instances. For a lookup by primary key, s.get() is shorter, and it checks the session's identity map before it sends any SQL.

python
from sqlalchemy import select

async def adults_and_first():
    async with Session() as s, s.begin():
        s.add(User(name="ada", age=36))          # sync: only stages the object
        res = await s.execute(select(User).where(User.age > 30))
        users = res.scalars().all()
        first = await s.get(User, 1)             # primary-key lookup
        return users, first

Assumes a mapped `User` class with `name` and `age` columns.

The result object has several ways to read rows, and the right one depends on how many rows you expect.

Result callReturns
scalars().all()A list of ORM objects
scalar_one()Exactly one object, otherwise an error
scalar_one_or_none()One object, or None when nothing matched
first()The first row, or None

Not every session method is a coroutine, and forgetting an await is an easy mistake. add() only places an object in the session, so it is an ordinary call. Anything that may touch the network has to be awaited.

MethodAwait it?
add()No, it is synchronous
flush(), commit()Yes
rollback(), refresh()Yes
delete(), execute()Yes
Common mistake

Writing s.commit() without await creates a coroutine that never runs, so nothing is saved. Python only warns that the coroutine was never awaited, and the warning is easy to miss.

Streaming, the greenlet bridge and run_sync

execute() fetches the whole result before it returns, which is a problem for a table with millions of rows. stream() is the alternative. It uses a server-side cursor, so rows arrive in batches and you consume them with async for. Memory stays flat no matter how large the result is.

python
async def export_events():
    async with Session() as s:
        result = await s.stream(select(Event))
        async for row in result:
            handle(row)          # one row at a time, never the full set

It is worth knowing how the ORM manages to be async at all. The ORM core is the same synchronous code that SQLAlchemy 1.x had. A greenlet bridges it to the async driver: when the sync code needs the network, the greenlet switches out, the driver awaits the I/O on the event loop, and execution switches back. This only works inside calls that go through the bridge, such as await s.execute().

If your code triggers database I/O outside that bridge, there is nothing to switch back to, and SQLAlchemy raises MissingGreenlet. The classic cause is reading an expired attribute or a lazy relationship such as user.posts straight from a plain attribute access. The cure is to load what you need up front or to await the load explicitly, which the N+1 section covers in depth.

Some SQLAlchemy APIs are synchronous only and cannot be awaited, for example schema creation and database inspection. For those, open an async connection and hand the sync function to run_sync, which runs it on the bridge with a sync-style connection.

python
async def create_tables():
    async with engine.begin() as conn:
        await conn.run_sync(Base.metadata.create_all)

DDL helpers like create_all expect a sync connection, so run_sync supplies one.

One session per task, and a clean shutdown

An AsyncSession tracks one unit of work: one transaction, one set of loaded objects, one connection. It is not safe for concurrent tasks. If two tasks run queries on the same session at the same time, their operations overlap inside a single transaction, and SQLAlchemy raises errors such as IllegalStateChangeError or the driver reports an operation already in progress. The rule is simple: share the engine and the sessionmaker, never the session.

Common mistake

Running asyncio.gather(*(s.get(User, i) for i in ids)) on one shared session. It is unsafe, and it would send N separate queries anyway. When you only need many rows, use one query such as select(User).where(User.id.in_(ids)).

When tasks really do need to run in parallel, give each task its own session. The async with inside the helper opens a session for that task alone and returns its connection to the pool when the task finishes.

python
import asyncio

async def load(user_id):
    async with Session() as s:
        return await s.get(User, user_id)

async def load_many(ids):
    return await asyncio.gather(*(load(i) for i in ids))

In a web framework the same idea becomes a dependency that opens one session per request. With FastAPI, a generator dependency yields the session and closes it after the response is sent.

python
async def get_session():
    async with Session() as s:
        yield s

@app.get("/users/{uid}")
async def read(uid: int, s: AsyncSession = Depends(get_session)):
    return await s.get(User, uid)

The last piece is shutdown. The engine keeps pooled connections open, so call await engine.dispose() when the application stops, or the database sees those connections drop abruptly. In FastAPI the natural place is the lifespan function: create tables on the way in, dispose on the way out.

python
from contextlib import asynccontextmanager

@asynccontextmanager
async def lifespan(app):
    async with engine.begin() as conn:
        await conn.run_sync(Base.metadata.create_all)
    yield
    await engine.dispose()

app = FastAPI(lifespan=lifespan)

When something goes wrong, the error name usually points to the cause.

ErrorLikely cause
MissingGreenletA lazy load or an expired attribute read outside the bridge
IllegalStateChangeErrorOne session shared by several tasks
coroutine was never awaitedA missing await on commit() or execute()
QueuePool limit reachedSessions opened but never closed
Remember

Create the engine and the sessionmaker once. Open one session per request or task with async with. Await everything except add(). Dispose of the engine on shutdown.

Part 6 · Transactions and Isolation Levels

Transaction blocks in asyncpg

A transaction groups several statements so they either all take effect or none do. In asyncpg you open one with conn.transaction() and use it as an async context manager. When the block exits cleanly, asyncpg sends COMMIT. If an exception escapes the block, it sends ROLLBACK and then re-raises the exception, so your error handling still sees it.

python
async def move_funds(conn, src, dst, amount):
    async with conn.transaction():
        await conn.execute(
            'UPDATE accounts SET bal = bal - $1 WHERE id = $2', amount, src)
        await conn.execute(
            'UPDATE accounts SET bal = bal + $1 WHERE id = $2', amount, dst)
    # COMMIT has run here; an exception above would have rolled both back

Both updates commit together or not at all

Transaction blocks can nest. A second conn.transaction() inside an open one does not start a new transaction, because Postgres has no such thing. asyncpg turns it into a SAVEPOINT instead. If the inner block fails, only the work since that savepoint is rolled back, and the outer transaction can carry on and commit what came before.

python
async def import_rows(conn, rows):
    async with conn.transaction():          # BEGIN
        for row in rows:
            try:
                async with conn.transaction():  # SAVEPOINT
                    await conn.execute(
                        'INSERT INTO items VALUES ($1, $2)', *row)
            except Exception:
                pass  # only this row's savepoint is rolled back
    # the good rows are committed together

One bad row does not undo the others

What a nested block does
  1. 1Outer block opensBEGIN
  2. 2Inner block opensSAVEPOINT
  3. 3Inner failsROLLBACK TO SAVEPOINT
  4. 4Outer exits cleanlyCOMMIT keeps earlier work

The transaction() call also takes options. isolation= picks the isolation level, readonly=True forbids writes, and deferrable=True is meaningful only together with serializable and read-only. A deferrable transaction waits until Postgres can hand it a snapshot that cannot be part of a serialization failure, then runs without risk of being aborted. That makes it a good fit for a long report.

python
async def big_report(conn, query):
    tx = conn.transaction(
        isolation='serializable',
        readonly=True,
        deferrable=True,
    )
    async with tx:
        return await conn.fetch(query)

A long read-only scan that will never fail with 40001

What each isolation level sees

An isolation level decides which other transactions' changes your transaction can see while it runs. Postgres offers three behaviours, and you choose between them per transaction or per engine.

The default is READ COMMITTED. Each statement sees the data that was committed before that statement began. Two identical SELECTs in the same transaction can therefore return different rows if someone commits in between. REPEATABLE READ takes a single snapshot when the transaction starts and uses it for every statement. In Postgres this blocks non-repeatable reads and also phantoms. SERIALIZABLE uses Serializable Snapshot Isolation (SSI), which tracks read and write dependencies between transactions and blocks every anomaly. The price is that Postgres may abort a transaction with a serialization failure.

You can also ask for READ UNCOMMITTED, but Postgres treats it exactly like READ COMMITTED. There are no dirty reads in Postgres at any level.

LevelDirty readNon-repeatable readPhantomSerialization anomaly
READ COMMITTEDnoyesyesyes
REPEATABLE READnonono*yes
SERIALIZABLEnononono

The asterisk matters: the SQL standard allows phantoms at REPEATABLE READ, but Postgres does not, because its snapshot covers the whole transaction. Do not assume the same on other databases.

LevelSnapshot takenCan fail with 40001?
READ COMMITTEDper statementno
REPEATABLE READonce per transactionyes, on write conflicts
SERIALIZABLEonce per transaction, plus SSI checksyes, and you must retry
Choosing a level
Rule of thumb

Most OLTP code is fine at READ COMMITTED. Use REPEATABLE READ for a consistent report, and SERIALIZABLE when a rule depends on several rows at once and you are willing to retry.

Retries, SQLAlchemy and row locks

SERIALIZABLE blocks every anomaly by aborting a transaction when it detects a dangerous pattern. The error is a SerializationFailure with SQLSTATE 40001. The database has already rolled the transaction back, so the fix is to rerun the whole transaction from BEGIN, not just the statement that raised. Back off a little between tries so the competing transactions can finish.

The loop below uses a fake transfer() that fails twice, so you can watch the retry logic without a database. With asyncpg you would catch asyncpg.SerializationError and wrap async with conn.transaction(isolation='serializable') inside the try.

python
import asyncio

class SerializationFailure(Exception):
    sqlstate = '40001'

attempts = 0

async def transfer():
    global attempts
    attempts += 1
    if attempts < 3:
        raise SerializationFailure('could not serialize access')
    return 'committed'

async def run_with_retry():
    for n in range(4):
        try:
            result = await transfer()
            print(f'attempt {n + 1}: {result}')
            return result
        except SerializationFailure as exc:
            print(f'attempt {n + 1}: {exc.sqlstate}, retrying')
            await asyncio.sleep(0.01 * 2**n)
    raise RuntimeError('gave up')

asyncio.run(run_with_retry())
output
attempt 1: 40001, retrying
attempt 2: 40001, retrying
attempt 3: committed

SQLAlchemy offers the same ideas. engine.execution_options(isolation_level='REPEATABLE READ') returns a copy of the engine, leaving the original unchanged, and every connection checked out through the copy uses that level. Inside a session, s.begin_nested() creates a savepoint. It needs an outer s.begin(), and if the nested block fails only the savepoint is rolled back.

python
async def audit_in_savepoint(engine, Session, audit):
    rr_engine = engine.execution_options(isolation_level='REPEATABLE READ')
    async with Session(bind=rr_engine) as s, s.begin():
        try:
            async with s.begin_nested():   # SAVEPOINT
                s.add(audit)
        except Exception:
            pass  # savepoint rolled back, outer work still commits

Level set on an engine copy, savepoint inside the outer transaction

Isolation levels govern snapshots; row locks let you claim rows explicitly. select(...).with_for_update() emits SELECT ... FOR UPDATE, which locks the chosen rows until the transaction ends, so no one else can change them meanwhile. For a job queue, add skip_locked=True: workers skip rows that another worker has already locked instead of waiting, so each worker claims a different job.

python
async def claim_job(Session, Job, select):
    async with Session() as s, s.begin():
        q = (select(Job)
             .where(Job.state == 'new')
             .limit(1)
             .with_for_update(skip_locked=True))
        job = (await s.execute(q)).scalar_one_or_none()
        if job is not None:
            job.state = 'running'
        return job

FOR UPDATE SKIP LOCKED through SQLAlchemy

The raw SQL version works the same way through asyncpg, as long as it runs inside a transaction so the lock lasts until COMMIT.

sql
SELECT id FROM jobs
WHERE state = 'new'
LIMIT 1
FOR UPDATE SKIP LOCKED;

Never await outside work inside a transaction

Because await hands control back to the event loop, it is easy to forget that the transaction stays open while your coroutine is suspended. If you await an HTTP call inside async with conn.transaction():, the database sits waiting for you. Any row locks you took are held the whole time, and the pooled connection is pinned, so other tasks queue for the pool.

Common mistake

Awaiting an HTTP call (or any external I/O) inside a transaction holds locks and a pooled connection for as long as the remote server takes. From the database side the session shows as idle in transaction.

python
async def bad(conn, http, url, upd, ins):
    async with conn.transaction():
        await conn.execute(upd)
        r = await http.get(url)      # locks held while we wait
        await conn.execute(ins, r)

async def good(conn, http, url, upd, ins):
    r = await http.get(url)          # outside the transaction
    async with conn.transaction():
        await conn.execute(upd)
        await conn.execute(ins, r)

Do external I/O first, then open a short transaction

To find leaks, query pg_stat_activity for sessions in that state. An old xact_start means a transaction has been open for a long time without running anything.

sql
SELECT pid, xact_start, query
FROM pg_stat_activity
WHERE state = 'idle in transaction';

As a safety net, set idle_in_transaction_session_timeout as a server setting, for example '10s'. Postgres then closes any session that stays idle inside a transaction longer than that, and the pool discards the dead connection and opens a fresh one on the next checkout.

Keep it short

Open the transaction late, close it early, and keep awaits on other services outside it. Nested blocks are savepoints, SERIALIZABLE means retry the whole transaction on 40001, and queue workers use FOR UPDATE SKIP LOCKED.

Part 7 · N+1 Queries and Lazy Loading Under asyncio

What N+1 is, and why AsyncSession refuses to hide it

The N+1 problem starts with one query that loads N parent rows, for example 100 users. You then read user.posts on each parent, and every read fires its own query. That is 1 + N round trips to Postgres, here 101. The code looks innocent, because the extra queries come from a plain attribute access. You only see them in the SQL log or in a slow endpoint.

One page of users, 101 queries
  1. 1SELECT users1 query, N rows
  2. 2posts WHERE user_id=1+1
  3. 3posts WHERE user_id=2+1
  4. 4... user_id=N+N in total

The example below fakes the database with a list of query strings, so you can count round trips without a server. The first loop is the N+1 shape: one query per user. The second asks for all the posts in one IN query, which is what an eager loader does for you.

python
import asyncio

USERS = {1: 'ada', 2: 'bo', 3: 'cy'}
POSTS = {1: ['a1', 'a2'], 2: ['b1'], 3: []}
queries = []

async def run(sql):
    queries.append(sql)
    await asyncio.sleep(0)

async def posts_for(uid):
    await run(f'SELECT * FROM posts WHERE user_id = {uid}')
    return POSTS[uid]

async def posts_for_many(ids):
    await run(f'SELECT * FROM posts WHERE user_id IN {tuple(ids)}')
    return {i: POSTS[i] for i in ids}

async def main():
    ids = list(USERS)
    await run('SELECT * FROM users')
    for i in ids:
        await posts_for(i)
    print('N+1:', len(queries), 'queries')
    queries.clear()
    await run('SELECT * FROM users')
    await posts_for_many(ids)
    print('IN query:', len(queries), 'queries')
    for q in queries:
        print(q)

asyncio.run(main())
output
N+1: 4 queries
IN query: 2 queries
SELECT * FROM users
SELECT * FROM posts WHERE user_id IN (1, 2, 3)

With the synchronous Session, a lazy load quietly runs the extra query when you touch the attribute. That is convenient, and it is also how N+1 hides. In an AsyncSession the attribute access is plain synchronous Python, so it has no way to await the query. Instead of silently querying, SQLAlchemy raises MissingGreenlet. The error is annoying, but it turns a hidden performance bug into an immediate failure that points at the line.

python
async def broken(s, User, select):
    users = (await s.scalars(select(User))).all()
    return users[0].posts   # implicit lazy load -> MissingGreenlet

Fragment: s is an AsyncSession and User is a mapped class.

Common mistake

Catching MissingGreenlet or turning it off is the wrong fix. The error means a relationship was never loaded, so load it on purpose, as the next pages show.

Eager loading: selectinload and joinedload

The cure for N+1 is to decide up front which relationships a query needs, and tell SQLAlchemy with a loader option. There are two main strategies. They differ in how many queries they run and in what shape the data comes back, so the right one depends on the relationship.

LoaderSQL it runsQueriesBest for
selectinload(User.posts)second query with WHERE user_id IN (...)2 totalone-to-many collections
joinedload(Post.author)a single LEFT OUTER JOIN1many-to-one
joinedload on a collectionJOIN with duplicated parent rows1only with .unique()
awaitable_attrs.postslazy SELECT, awaited+1a deliberate single load
raiseload('*')noneerrortests

selectinload is the default choice for collections. It runs your main query first, then a second query that fetches all the children whose foreign key is in the list of parent ids it just loaded. Two queries cover any number of parents, and the parent rows are never repeated in the result.

python
async def users_with_posts(s, User, select, selectinload):
    stmt = select(User).options(selectinload(User.posts))
    users = (await s.scalars(stmt)).all()
    return [(u.name, len(u.posts)) for u in users]   # no extra queries

Fragment: the posts are already loaded, so reading u.posts is safe.

joinedload folds the related rows into the main query with a LEFT OUTER JOIN. That fits a many-to-one link such as a post's author, because each post row carries exactly one author and nothing is duplicated. If you use it on a collection, each parent appears once per child row, and SQLAlchemy refuses to hand you those duplicates silently. You must call result.unique() before reading the rows.

python
async def posts_with_authors(s, Post, select, joinedload):
    stmt = select(Post).options(joinedload(Post.author))
    posts = (await s.scalars(stmt)).all()      # many-to-one: fine
    return [p.author.name for p in posts]

async def users_joined(s, User, select, joinedload):
    stmt = select(User).options(joinedload(User.posts))
    result = await s.scalars(stmt)
    return result.unique().all()               # required for collections
Common mistake

Using joinedload on a collection without calling .unique() raises an error, and even when you work around it, the JOIN repeats every parent row once per child. Use selectinload for collections.

Failing loudly, loading on purpose, and surviving commit

Eager loading only helps if you remember it everywhere. To catch the places you forgot, make unplanned lazy loads fail on every code path. Declare the relationship with relationship(lazy="raise"), or add options(raiseload('*')) to a query. Any access to a relationship you did not load explicitly then raises straight away. This is especially useful in tests, because an endpoint that would have been N+1 in production fails in the test suite instead.

python
async def strict_query(s, User, select, raiseload, selectinload):
    stmt = select(User).options(
        selectinload(User.posts),   # planned
        raiseload('*'),             # everything else raises
    )
    return (await s.scalars(stmt)).all()

Sometimes you do want a lazy load, for example to fetch posts for one user you already hold. Mix the AsyncAttrs mixin into your declarative base (or a model class), and every relationship gets an awaitable twin under awaitable_attrs. Awaiting it runs the query visibly, so nothing is hidden. If the relationship is already loaded and you want fresh data, await s.refresh(user, ["posts"]) reloads just that relationship.

python
async def load_on_purpose(s, User):
    user = await s.get(User, 1)
    posts = await user.awaitable_attrs.posts   # explicit, awaited query
    await s.refresh(user, ['posts'])           # reload one relationship
    return posts

Fragment: User must include the AsyncAttrs mixin.

Commit has a lazy-loading trap of its own. By default a session expires every object on commit, so the next attribute read tries to refresh the row from the database. That read is plain synchronous access again, so it raises MissingGreenlet. Create the factory with async_sessionmaker(engine, expire_on_commit=False) and the attributes you already loaded stay readable after commit. You trade away automatic freshness, so reload with refresh() when you need current values.

SituationWithout the fixWith the fix
Read user.name after commitrefresh needed, MissingGreenletexpire_on_commit=False keeps the value
Read an unloaded user.postsMissingGreenletawait user.awaitable_attrs.posts
Want fresh posts for one userstale or unloadedawait s.refresh(user, ['posts'])
Hunt hidden lazy loadserrors only in production pathslazy='raise' fails in tests
Common mistake

Leaving the default expire_on_commit=True and then reading attributes after commit makes the next access trigger a hidden refresh, which raises MissingGreenlet.

Batching at the asyncpg level, and a tempting gather trap

Plain asyncpg has no ORM and no lazy loading, but it has the same N+1 shape: a loop that awaits fetchrow() once per id. The fix is to send the whole list of ids in one query. Postgres accepts an array parameter, so you write WHERE id = ANY($1::int[]) and pass an ordinary Python list. One round trip returns every matching row.

python
async def users_by_ids(conn, ids):
    rows = await conn.fetch(
        'SELECT * FROM users WHERE id = ANY($1::int[])',
        ids,                      # a Python list, sent as one array
    )
    return {r['id']: r for r in rows}

Fragment: conn is an asyncpg connection.

A second trap looks modern rather than naive. It is tempting to avoid the loop by firing all the lookups at once with asyncio.gather. On a single AsyncSession that does not help, and it is not safe.

python
async def tempting_but_wrong(s, User, ids):
    return await asyncio.gather(*(s.get(User, i) for i in ids))

Fragment: N queries on one shared session.

This still sends N separate queries. Worse, one session wraps one connection and one transaction, and it is not task-safe, so running many operations on it at the same time can raise IllegalStateChangeError or interleave state. The right answer is the same as before: ask the database for all the ids in a single IN query.

python
async def one_in_query(s, User, select, ids):
    stmt = select(User).where(User.id.in_(ids))
    return (await s.scalars(stmt)).all()      # 1 query, safe on one session
Rule of thumb

Whenever you see a loop or a gather that awaits one query per item, look for a way to send the whole set in one query. Use selectinload for collections, joinedload for many-to-one, and ANY or IN for raw id lists.

Part 8 · Connection Pool Sizing

Why a pool is small and finite

Every Postgres connection is a separate backend process on the server, and each one costs several MB of RAM plus a share of CPU scheduling. That is why the server caps them: max_connections defaults to 100, and 3 of those slots are reserved for superusers, so your applications can really count on 97. Past that point the database refuses new connections, and well before it, extra connections only add memory pressure and contention.

Your side of the connection is managed by a pool. SQLAlchemy's QueuePool has four settings that matter, and their defaults are worth knowing by heart because the default error message quotes them.

SettingDefaultMeaning
pool_size5Connections kept open permanently
max_overflow10Extra connections allowed during bursts, closed again afterwards
pool_timeout30 sHow long a task waits for a free connection before raising an error
pool_recycle-1Connections are never replaced because of age

So by default one process can hold at most 5 + 10 = 15 connections, and the 16th caller waits up to 30 seconds.

In asyncio the pool is a semaphore

In a threaded server you might need roughly one connection per thread. In asyncio the situation is different: thousands of coroutines can be in flight, but a coroutine only needs a connection while a query or transaction is actually running. The pool behaves like a semaphore: every coroutine that wants a connection waits in line, and waiting costs almost nothing because the event loop simply serves other tasks meanwhile. A small pool is therefore normal, not a sign of a misconfiguration.

Working out the right number

Start from demand. Little's law says the number of connections in use at any moment is about the arrival rate multiplied by the average time each request holds a connection. Use the time the connection is checked out, not the full request time: a request that spends 200 ms calling another service but only 20 ms inside a transaction holds a connection for 20 ms.

python
rps = 500          # requests per second
hold = 0.02        # seconds a connection is held per request
needed = rps * hold
print(f"connections needed: {round(needed)}")

500 requests per second at 20 ms each keep about 10 connections busy.

output
connections needed: 10

A second, independent check comes from the PostgreSQL wiki. As a starting heuristic, the useful number of connections is about (CPU cores × 2) + effective spindles. A database server with 8 cores and one SSD lands near 17. This surprises people, but a bigger pool often makes things slower: more queries compete for the same CPU and disk, and every extra backend adds context switching and lock contention. If Little's law gives 10 and the heuristic gives 17, a pool in that range is plenty.

The budget across your whole fleet

The pool belongs to one process, but max_connections belongs to the whole database. The budget is: app instances × workers per instance × (pool_size + max_overflow) must stay below max_connections minus the headroom you keep for admin work such as psql sessions, migrations and cron jobs. Note that max_overflow counts, because a burst can use it on every process at once.

python
def total(instances, workers, pool_size, max_overflow):
    return instances * workers * (pool_size + max_overflow)

limit = 100 - 3 - 7   # max_connections, superuser slots, admin headroom
for name, args in [("3 x 4, 5+10", (3, 4, 5, 10)),
                   ("2 x 4, 3+2", (2, 4, 3, 2))]:
    used = total(*args)
    print(f"{name}: {used} of {limit} -> fits: {used < limit}")

The same function checks the default pool and a shrunken one.

output
3 x 4, 5+10: 180 of 90 -> fits: False
2 x 4, 3+2: 40 of 90 -> fits: True

Three instances with four workers each and the default pool could open 180 connections against 90 available. The fix is to shrink the per-process pool: with 3 + 2 per process, two instances of four workers need only 2 × 4 × (3 + 2) = 40 connections. If you autoscale, budget for the maximum instance count, and when trimming, reduce max_overflow before pool_size.

DeploymentPer processTotalFits within 90?
1 instance × 1 worker5 + 1015yes
2 instances × 4 workers5 + 10120no
2 instances × 4 workers3 + 240yes
3 instances × 4 workers3 + 260yes
Any, with PgBouncer and NullPool0 in the appset by the pooleryes

Keeping connections healthy

A pool also keeps connections open for a long time, and a connection can die while it sits idle: the database restarts, a firewall drops idle flows, or a load balancer or proxy silently closes the link. The pool does not notice until you try to use it, and then your query fails. Three settings address this, each with a different cost.

SettingProtects againstCost
pool_pre_ping=TrueDatabase restarts, firewall idle-killsOne extra round trip on every checkout
pool_recycle=1800Load balancers or proxies that drop idle links without telling anyoneA reconnect every 30 minutes per connection
poolclass=NullPoolDouble pooling behind PgBouncer or on serverlessA new connection on every checkout

Pre-ping sends a lightweight check before handing out a connection and replaces it if the check fails, so a restart of the database costs you a moment instead of a burst of errors. Recycling replaces connections before they get old enough to be dropped, so set the value below the shortest idle timeout on the path to the database; 1800 seconds is a common choice. Together they give a robust configuration for a direct connection:

text
engine = create_async_engine(
    url,
    pool_size=5,
    max_overflow=10,
    pool_timeout=30,
    pool_pre_ping=True,
    pool_recycle=1800,
)

Fragment: needs SQLAlchemy and a real database URL.

When an external pooler such as PgBouncer sits in front of Postgres, or when you run on serverless where processes come and go, an in-app pool just stacks a second pool on top of the first and wastes slots. Use NullPool, so each checkout opens a connection and each release closes it, and let the external pooler own the pooling. In PgBouncer's transaction mode also turn off asyncpg's statement caches, as the prepared statements section explained.

text
engine = create_async_engine(
    url,
    poolclass=NullPool,
    connect_args={"statement_cache_size": 0},
)

Fragment: the pooler caps the connections, not the application.

Diagnosing pool trouble

When the pool runs dry, a task waits pool_timeout seconds and then fails with an error that quotes your own settings. This is the signature symptom of an undersized pool or of sessions that never give their connection back.

output
TimeoutError: QueuePool limit of size 5 overflow 10 reached, connection timed out, timeout 30

Because both causes produce the same message, check for leaks first. If every session is opened with async with, and no HTTP call or other slow external work sits inside a transaction, the leak is unlikely and the pool really is too small for the load. The chart below walks through the usual order of questions.

Pool timeouts: where to look

Two tools help you see what is going on. In Postgres, SELECT state, count(*) FROM pg_stat_activity GROUP BY state; shows how many connections are active, idle, or stuck as idle in transaction; a growing last group means leaked sessions or held transactions. In the application, engine.pool.status() reports how many connections are checked out and how much overflow is in use. Before you raise pool_size, cut the hold time: fewer queries per checkout and no external calls inside a transaction reduce the need under Little's law.

Common mistake

Sizing for one process and forgetting to multiply by instances and workers, so the fleet quietly asks for more than max_connections allows.

Common mistake

Setting a pool of 100 'for speed'. The database has to run all those queries on the same cores, so each one gets slower through contention.

Common mistake

Keeping an in-app pool behind PgBouncer instead of using NullPool, or opening sessions without async with so connections are never returned.

Why can 5 connections serve thousands of concurrent requests?
  • A request only holds a connection while a query or transaction runs, typically a few milliseconds, and waits in line for it without blocking the event loop.
  • Little's law shows the real need: 500 requests per second at 20 ms each keep only about 10 connections busy, so a pool of 5 + 10 covers it.
Remember

Size the pool from hold time and database cores, then check the product of instances, workers and pool slots against max_connections minus headroom. Add pool_pre_ping and pool_recycle for health, and use NullPool when something else owns the pooling.

Part 9 · aiosqlite and redis.asyncio

aiosqlite: SQLite without blocking the loop

The standard sqlite3 module is synchronous, so calling it directly inside async def freezes the event loop for as long as the query runs. aiosqlite fixes this without rewriting SQLite. Each connection gets its own dedicated background thread, and that thread runs the real stdlib sqlite3 code. The async methods you call, such as execute and commit, only put a request on that thread's queue and then await the answer.

The sketch below is a stripped-down version of the same idea, built only from the standard library so you can run it. The worker thread owns the sqlite3 connection, the coroutine hands it work through a queue, and a future carries the result back to the loop. The real library adds cursors, transactions and error handling on top of this shape.

python
import asyncio, queue, sqlite3, threading

class TinyAsyncDB:
    def __init__(self):
        self._jobs = queue.Queue()
        self.last_thread = None
        self._thread = threading.Thread(target=self._run, name='sqlite-worker', daemon=True)
        self._thread.start()

    def _run(self):
        conn = sqlite3.connect(':memory:')
        while (job := self._jobs.get()) is not None:
            fut, loop, sql, params = job
            try:
                rows = conn.execute(sql, params).fetchall()
                conn.commit()
                self.last_thread = threading.current_thread().name
                loop.call_soon_threadsafe(fut.set_result, rows)
            except Exception as exc:
                loop.call_soon_threadsafe(fut.set_exception, exc)
        conn.close()

    async def execute(self, sql, params=()):
        loop = asyncio.get_running_loop()
        fut = loop.create_future()
        self._jobs.put((fut, loop, sql, params))
        return await fut

    def close(self):
        self._jobs.put(None)
        self._thread.join()

async def main():
    db = TinyAsyncDB()
    await db.execute('CREATE TABLE t (n INTEGER)')
    await db.execute('INSERT INTO t VALUES (?)', (7,))
    rows = await db.execute('SELECT n FROM t')
    print('rows:', rows)
    print('event loop thread:', threading.current_thread().name)
    print('sqlite thread:', db.last_thread)
    db.close()

asyncio.run(main())

A toy version of what aiosqlite does for every connection

output
rows: [(7,)]
event loop thread: MainThread
sqlite thread: sqlite-worker

The loop never touches SQLite itself; it only awaits a future. That is why other coroutines keep running while a query executes. It also means a single aiosqlite connection still runs one statement at a time, because its worker thread is a single queue.

Connecting, querying and committing

In real code you open the connection with async with aiosqlite.connect("app.db") as db:, which starts the thread on entry and stops it on exit. Run statements with await db.execute(sql, params), always passing values as parameters rather than formatting them into the SQL. Writes are not durable until you call await db.commit(); leaving the block without committing discards them. Setting db.row_factory = aiosqlite.Row lets you read columns by name, as in row["name"], instead of by position.

python
async def list_adults(min_age):
    import aiosqlite

    async with aiosqlite.connect('app.db') as db:
        db.row_factory = aiosqlite.Row
        await db.execute('INSERT INTO users (name, age) VALUES (?, ?)', ('ada', 36))
        await db.commit()  # nothing is saved without this

        async with db.execute('SELECT name, age FROM users WHERE age >= ?', (min_age,)) as cur:
            rows = await cur.fetchall()
        return [(r['name'], r['age']) for r in rows]

Fragment: call it from a running event loop, for example with asyncio.run(list_adults(18))

Common mistake

Calling plain sqlite3 from a coroutine because "it is only a local file". Even a fast query blocks every other task while it runs, and a slow one stalls the whole service. Use aiosqlite, or wrap the call in asyncio.to_thread().

One writer at a time: WAL, busy_timeout and SQLAlchemy

SQLite lets only one writer hold the database at a time. When a second connection tries to write while the first is mid-transaction, it waits for the lock. Python's sqlite3.connect has a default timeout of 5.0 seconds, and that wait is what acts as a busy timeout. You only see database is locked once that wait runs out, or immediately if the timeout is set to 0. Under real load, 5 seconds of waiting is often not enough, and the error shows up as a flaky failure.

Two settings address this. PRAGMA journal_mode=WAL switches SQLite to write-ahead logging, so readers no longer block the writer and the writer no longer blocks readers. The setting is stored in the database file, so it persists. PRAGMA busy_timeout=5000 (or the timeout argument of connect) tunes how many milliseconds a connection waits for the write lock before giving up. This one is per connection, so set it every time you open one.

The demo below uses plain sqlite3 so it can run anywhere. One connection starts a write and holds it open, while a second connection with a short 0.1 s timeout tries to write. The second connection gives up with database is locked, yet it can still read, because WAL keeps readers out of the writer's way.

python
import os, sqlite3, tempfile

path = os.path.join(tempfile.mkdtemp(), 'app.db')
a = sqlite3.connect(path, isolation_level=None)
print('journal mode:', a.execute('PRAGMA journal_mode=WAL').fetchone()[0])
a.execute('CREATE TABLE t (n INTEGER)')
a.execute('BEGIN IMMEDIATE')
a.execute('INSERT INTO t VALUES (1)')

b = sqlite3.connect(path, timeout=0.1, isolation_level=None)
print('busy_timeout ms:', b.execute('PRAGMA busy_timeout').fetchone()[0])
try:
    b.execute('BEGIN IMMEDIATE')
except sqlite3.OperationalError as exc:
    print('second writer:', exc)
print('reader sees:', b.execute('SELECT count(*) FROM t').fetchone()[0])

a.execute('COMMIT')
print('after commit:', b.execute('SELECT count(*) FROM t').fetchone()[0])
a.close()
b.close()
output
journal mode: wal
busy_timeout ms: 100
second writer: database is locked
reader sees: 0
after commit: 1

With aiosqlite you run the same two pragmas right after opening the connection. Keep write transactions short, because every millisecond you hold the lock is a millisecond other writers spend waiting.

python
async def open_db():
    import aiosqlite

    db = await aiosqlite.connect('app.db')
    await db.execute('PRAGMA journal_mode=WAL')
    await db.execute('PRAGMA busy_timeout=5000')
    return db

Fragment: close the returned connection with await db.close()

Common mistake

Leaving WAL and the timeout at their defaults and then blaming asyncio for database is locked. The loop is fine; two writers simply collided, and the second one stopped waiting too soon.

SQLAlchemy over SQLite for tests and local development

If your app uses SQLAlchemy 2.0, the URL sqlite+aiosqlite:///app.db runs the same AsyncSession code on a local file with no server to install. That makes it handy for unit tests and for a first run on a laptop. Treat it as a stand-in, though: SQLite and Postgres behave differently in ways that matter, and a test that passes on one can still fail on the other.

python
def make_test_engine():
    from sqlalchemy.ext.asyncio import create_async_engine

    return create_async_engine('sqlite+aiosqlite:///app.db')

Swap the URL for postgresql+asyncpg://... in production

Aspectsqlite+aiosqlitepostgresql+asyncpg
WritersOne at a timeMany at once (MVCC)
TypesLoose affinity, a column accepts almost anythingStrict, a wrong type is an error
LockingWhole database fileRow level
Best forTests and local developmentProduction

The practical consequence is that lock-related bugs, such as deadlocks, SELECT ... FOR UPDATE behaviour and isolation levels, will not show up on SQLite. Run at least part of your test suite against a real Postgres before you trust it.

redis.asyncio: clients, pipelines and pools

The old standalone aioredis library was merged into redis-py in version 4.2, and it now lives there as redis.asyncio. You get the same commands as the sync client, but each one is a coroutine. The usual way to create a client is redis.asyncio.Redis.from_url(url, decode_responses=True). With decode_responses=True, replies come back as str instead of bytes, which is what you want in most application code.

python
async def cache_demo():
    import redis.asyncio as redis

    r = redis.Redis.from_url('redis://localhost:6379/0', decode_responses=True)
    await r.set('k', 'v', ex=60)   # expires after 60 seconds
    value = await r.get('k')       # 'v' as a str
    await r.aclose()               # redis-py 5.0.1 and later
    return value

Fragment: run it inside an event loop with a Redis server available

The ex=60 argument gives the key a 60 second lifetime, and aclose() releases the client's connections. On redis-py versions before 5.0.1 the method was called close(), so check which one your installed version provides.

Pipelines batch round trips

Every awaited command costs a network round trip. When you have several commands to send, a pipeline queues them on the client and sends them together. With transaction=True the group is also wrapped in MULTI/EXEC, so the server runs the commands as one unit. Inside the block, the commands are queued without being awaited; only execute() is awaited, and it returns the list of replies.

A client is cheap to use but not free to create, because each one owns connections. Build one ConnectionPool for the whole application with ConnectionPool.from_url(url, max_connections=50), wrap it in a single Redis client, and share that client everywhere. Creating a new client for every request opens a new set of sockets each time and can overwhelm the Redis server under load.

python
async def count_hit():
    from redis.asyncio import ConnectionPool, Redis

    pool = ConnectionPool.from_url('redis://localhost:6379/0', max_connections=50, decode_responses=True)
    r = Redis(connection_pool=pool)   # create once at startup, share it

    async with r.pipeline(transaction=True) as p:
        p.incr('hits')                # queued, not sent yet
        p.expire('hits', 60)          # queued, not sent yet
        results = await p.execute()   # one round trip, MULTI/EXEC

    await r.aclose()                  # at shutdown
    await pool.disconnect()           # a pool you passed in is not closed for you
    return results

Fragment: in a real app the pool and client live for the whole process, not for one call

Because you supplied the pool yourself, aclose() does not tear it down. Call pool.disconnect() as well during application shutdown, for example in the lifespan handler.

Common mistakes

Creating a Redis client per request causes a connection storm. Skipping aclose() at shutdown leaks sockets. Awaiting each incr and expire separately instead of using a pipeline wastes round trips.

Pub/sub, optimistic locking and picking a tool

Pub/sub

Redis can push messages to subscribers, so a consumer does not have to poll. Call r.pubsub() to get a subscription object, await ps.subscribe("ch") to join a channel, and then iterate with async for msg in ps.listen(). The loop runs for as long as the subscription lives. The first item you receive is not data but a confirmation with type set to subscribe, so filter on msg["type"] == "message" before reading msg["data"]. If you would rather poll with a deadline, get_message(ignore_subscribe_messages=True, timeout=1.0) returns the next real message or None.

python
async def listen(r):
    ps = r.pubsub()
    await ps.subscribe('ch')
    try:
        async for msg in ps.listen():
            if msg['type'] == 'message':   # skip the subscribe confirmation
                print(msg['data'])
    finally:
        await ps.aclose()

Fragment: r is the shared client from the previous page. Run it as its own task, because it never ends on its own

Optimistic locking with WATCH

Suppose you read a balance, subtract from it and write it back. Another client could change the balance between your read and your write, and you would overwrite its update. Optimistic locking solves this without holding a lock. You watch the key, read it, then queue your write between multi() and execute(). If anyone else touched the watched key in the meantime, Redis refuses the transaction and execute() raises WatchError. Your job is then to retry the whole block from the read.

WATCH, MULTI and EXEC with retry
python
async def spend(r, amount):
    from redis.exceptions import WatchError

    async with r.pipeline() as p:
        for _ in range(5):                 # cap the retries
            try:
                await p.watch('bal')
                bal = int(await p.get('bal'))
                p.multi()                  # from here commands are queued
                p.set('bal', bal - amount)
                await p.execute()          # raises WatchError if bal changed
                return bal - amount
            except WatchError:
                continue                   # someone else won; redo the read
    raise RuntimeError('balance kept changing')

Fragment: r is the shared client

Ignoring WatchError means a lost update, which is the very bug the check exists to prevent. Under heavy contention on one key, though, retries pile up and most attempts fail. In that case move the logic into a Lua script. Redis runs a script atomically on the server, so there is no watch, no retry loop and only one round trip.

python
async def spend_atomic(r, amount):
    script = r.register_script("""
    local bal = tonumber(redis.call('GET', KEYS[1]) or '0')
    if bal < tonumber(ARGV[1]) then return -1 end
    return redis.call('DECRBY', KEYS[1], ARGV[1])
    """)
    return await script(keys=['bal'], args=[amount])   # -1 means insufficient funds

Fragment: the check and the update happen in one atomic step

Which Redis tool for which job

NeedUseWhy
Send many commands at oncepipeline()One round trip
All-or-nothing groupMULTI/EXECQueued and applied together
Read, then write based on itWATCH plus retryOptimistic lock
Complex atomic step under contentionLua scriptRuns on the server
Fan out events to listenerspubsub()Pushed, no polling
What to remember

Both libraries keep the event loop free, but in different ways: aiosqlite hands each call to a background thread, while redis.asyncio waits on the network. For SQLite, turn on WAL and a busy timeout; for Redis, share one pool and close it at shutdown.

Part 10 · Timeouts and Cancellation Safety

Deadlines on Every Layer

A database call can hang for many reasons: a lock that never frees, a query with a bad plan, or a network that silently dropped. Without a deadline the coroutine awaits forever, and it holds a pooled connection the whole time. The cure is a deadline at every layer, so whichever one is closest to the problem gives up first.

The client-side deadline

The simplest deadline lives in your own code. On Python 3.11+ wrap the awaits in async with asyncio.timeout(5):. If the block runs longer than five seconds, asyncio cancels the current await and the block raises the builtin TimeoutError. The example uses a fake query so it runs anywhere.

python
import asyncio

async def slow_query():
    await asyncio.sleep(1)
    return 'rows'

async def main():
    try:
        async with asyncio.timeout(0.1):
            print(await slow_query())
    except TimeoutError:
        print('gave up after 0.1 s')

asyncio.run(main())
output
gave up after 0.1 s

On older Pythons use await asyncio.wait_for(coro, 5), which wraps a single awaitable. Before 3.11 it raises asyncio.TimeoutError, which is a different class from the builtin TimeoutError. The two only became the same class in 3.11, so code that must run on older versions should catch asyncio.TimeoutError. The timeout() block is usually nicer because it covers several awaits with one shared budget.

Timeouts built into asyncpg

asyncpg has its own timeouts, so you often do not need to wrap each call yourself. Pass command_timeout= to asyncpg.connect() or create_pool() to set a default for every command on those connections. Pass timeout= to a single fetch() or execute() call to override it. When the limit passes, asyncpg does not just abandon the call. It sends a Postgres cancel request, so the server stops working on the query too.

text
pool = await asyncpg.create_pool(dsn, command_timeout=10)

# per-call override for one query
rows = await pool.fetch(q, timeout=2)

Pool-wide default of 10 s, with a tighter 2 s for one call.

Guards that live on the server

Client deadlines only help while your process is alive and healthy. Postgres can enforce limits itself, and they keep working if the client crashes or stalls. You set them as server settings when each connection opens, so every pooled connection carries them.

Server settingWhat it stops
statement_timeoutA runaway query
lock_timeoutA long wait for a lock
idle_in_transaction_session_timeoutA transaction left open with nothing running

With asyncpg pass them as server_settings= to create_pool(). With SQLAlchemy put the same dictionary under connect_args={'server_settings': {...}} in create_async_engine().

text
engine = create_async_engine(
    url,
    connect_args={
        'server_settings': {
            'statement_timeout': '5s',
            'lock_timeout': '2s',
            'idle_in_transaction_session_timeout': '10s',
        }
    },
)
Order the deadlines

Set statement_timeout a little below the client timeout. The server error then arrives first, with a clear message and a connection in a clean state, instead of the client cancelling mid-flight.

What Happens When a Task Is Cancelled

A timeout, a closed client connection or a failed sibling task all end the same way: asyncio cancels the task. The task is not killed on the spot. Instead a CancelledError is raised at the await where it is currently suspended, and from there it unwinds like any exception, running finally blocks and context-manager exits on the way out.

This is what makes database code mostly safe by default. If the task is inside async with conn.transaction():, the exception leaves the block, so the transaction issues a ROLLBACK. The pool then either resets the connection or discards it if its state looks doubtful, so the next user gets a clean one.

How a cancel unwinds
  1. 1Cancel arrivestask.cancel() or a timeout fires
  2. 2CancelledErrorraised at the current await
  3. 3ROLLBACKconn.transaction() exits with an error
  4. 4Pool cleans upconnection is reset or discarded
  5. 5Re-raisethe caller sees the cancel

The example below imitates a transaction with a tiny context manager so you can watch the order of events. Notice that the rollback happens before our own cleanup message, and that the cancel still reaches the caller.

python
import asyncio

class FakeTransaction:
    async def __aenter__(self):
        print('BEGIN')
        return self
    async def __aexit__(self, exc_type, exc, tb):
        print('ROLLBACK' if exc_type else 'COMMIT')
        return False

async def worker():
    try:
        async with FakeTransaction():
            print('UPDATE sent')
            await asyncio.sleep(10)
    except asyncio.CancelledError:
        print('cleaning up')
        raise

async def main():
    task = asyncio.create_task(worker())
    await asyncio.sleep(0.05)
    task.cancel()
    try:
        await task
    except asyncio.CancelledError:
        print('caller sees the cancel')

asyncio.run(main())
output
BEGIN
UPDATE sent
ROLLBACK
cleaning up
caller sees the cancel

Never swallow the cancel

Since Python 3.8, CancelledError inherits from BaseException, not Exception. A plain except Exception: lets it through, which is the intent. But except BaseException: pass or a bare except: catches it and discards it, and the task carries on as if nothing happened. The timeout block above then never fires correctly, and the task cannot be stopped at shutdown. The rule is simple: do your cleanup, then re-raise.

Common mistake: swallowing the cancel

except BaseException: pass hides CancelledError. A task that ignores a cancel keeps its connection, ignores asyncio.timeout() and blocks shutdown. Catch Exception for ordinary errors, and when you must catch the cancel, finish cleanup and raise.

Half-Read Connections, Shield and TaskGroup

A cancel in the middle of a command

A database connection is a conversation: you send a command and then read its reply from the socket. If a cancel lands after the command was sent but before the reply was fully read, the connection is left half-read. The next command on it may read the leftover bytes of the previous answer. Careful drivers detect this and close the connection, or run a reset before reuse. Reusing such a connection yourself is dangerous.

This is not hypothetical. In March 2023 redis-py had exactly this class of bug, tracked as CVE-2023-28858 and CVE-2023-28859. A cancelled command on a pooled async connection returned the connection to the pool with an unread reply, and a later request, possibly from a different user, received that reply. The fix is in the driver, but the lesson is for you too: keep drivers updated, never push a connection back into a pool by hand after a cancel, and if you are unsure of its state, close it.

Common mistake: reusing a cancelled connection

A connection cancelled mid-read may still hold stale replies. Let the pool release or discard it through async with pool.acquire(), rather than catching the cancel and carrying on with the same connection.

Protecting critical cleanup with shield

Sometimes a short step must finish even if the caller is cancelled, such as a final commit or releasing a connection. asyncio.shield() wraps the awaitable in its own task. If the outer task is cancelled, the inner task keeps running. There is a catch: the outer await still raises CancelledError straight away, so shield protects the work, not the caller. Use it for short cleanup only, never to make a whole job uncancellable.

python
import asyncio

async def final_commit():
    await asyncio.sleep(0.1)
    print('commit finished')

async def worker():
    try:
        await asyncio.shield(final_commit())
    except asyncio.CancelledError:
        print('worker cancelled')
        raise

async def main():
    task = asyncio.create_task(worker())
    await asyncio.sleep(0.02)
    task.cancel()
    try:
        await task
    except asyncio.CancelledError:
        print('outer await raised')
    await asyncio.sleep(0.2)
    print('done')

asyncio.run(main())
output
worker cancelled
outer await raised
commit finished
done

Cancellation inside a TaskGroup

A TaskGroup cancels all the remaining tasks as soon as one of them fails, then raises the errors together as an exception group. That is a good fit for database work, because a failed query should not leave its siblings running. It also means every sibling must be cancel-safe, and it means they must not share a connection or an AsyncSession. One connection serves one command at a time, and a cancel in one task could leave the shared connection half-read for the others. Give each task its own pool.acquire() or its own session.

python
import asyncio

async def query(name, delay, fail=False):
    try:
        await asyncio.sleep(delay)
        if fail:
            raise RuntimeError(f'{name} failed')
        print(name, 'ok')
    except asyncio.CancelledError:
        print(name, 'cancelled')
        raise

async def main():
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(query('A', 0.05, fail=True))
            tg.create_task(query('B', 1))
    except* RuntimeError as eg:
        print('group raised:', eg.exceptions[0])

asyncio.run(main())
output
B cancelled
group raised: A failed
Common mistake: one connection for the whole group

Passing one conn or one AsyncSession to every task in a TaskGroup breaks as soon as one task fails and the others are cancelled mid-command. Acquire inside each task.

Retrying Safely and Choosing a Timeout Layer

Once a deadline fires, the temptation is to just try again. That is only safe when running the work twice has the same effect as running it once. A transaction that failed with a serialization error was rolled back completely, so nothing happened and a retry is fine. A timed-out INSERT is different: the client gave up, but the server may already have committed the row. A blind retry then creates a duplicate.

What happenedRetry?
SerializationFailure (SQLSTATE 40001)Yes, rerun the whole transaction with backoff
Connection lost before the workYes, if the work is idempotent
Timeout on an INSERTNo, it may have committed. Use an idempotency key or check first
Deadlock or lock timeoutYes, with backoff, if nothing outside the database changed
Should I retry?

The loop below retries a fake serialization failure. It sleeps a little longer after each failure, which is called exponential backoff, so a crowd of clients does not hammer the database in lockstep. It reruns the whole unit of work each time, never a single statement from the middle of it. In real code add random jitter to the delay and stop after about three tries.

python
import asyncio

class SerializationFailure(Exception):
    pass

attempts = 0

async def transfer():
    global attempts
    attempts += 1
    if attempts < 3:
        raise SerializationFailure('40001')
    return 'transferred'

async def main():
    for n in range(4):
        try:
            print(await transfer())
            break
        except SerializationFailure:
            delay = 0.01 * 2**n
            print(f'attempt {n + 1} failed, sleeping {delay:.2f}s')
            await asyncio.sleep(delay)

asyncio.run(main())
output
attempt 1 failed, sleeping 0.01s
attempt 2 failed, sleeping 0.02s
transferred
Common mistake: re-running an INSERT after a timeout

The timeout only tells you the client stopped waiting, not that the server did nothing. Retrying a plain INSERT can create the row twice. Use a unique idempotency key with ON CONFLICT DO NOTHING, or look the row up before you retry.

Which layer fires when

LayerSet withOn expiry
asyncioasyncio.timeout(5)TimeoutError, the task is cancelled
asyncpgcommand_timeout= or timeout=A Postgres cancel request is sent
Postgresstatement_timeoutThe query is aborted by the server
Postgresidle_in_transaction_session_timeoutThe session is closed
Putting it together

Put a deadline on both sides: a client timeout plus server guards. Clean up on cancel and then re-raise CancelledError. Use shield() only for a short final commit or release. Give each concurrent task its own connection or session. Retry only idempotent work, with backoff.

Part 11 · Cheatsheet: Async Databases in Python

Setup: asyncpg and SQLAlchemy

This page collects the calls you reach for most often. Code that needs a live database is shown as a fragment: the call site only, with pool, dsn, uid and the other names assumed to exist. The runnable example is on the last page.

asyncpg: pool, connection, query

Create one pool per process at startup. Borrow a connection with async with pool.acquire(), which hands it back on exit. Then call one of four methods. Values go in as positional $1, $2 placeholders and travel separately from the SQL text, so %s and ? do not work.

text
pool = await asyncpg.create_pool(dsn)

async with pool.acquire() as c:
    n = await c.fetchval('SELECT count(*) FROM users WHERE org = $1', org)
    row = await c.fetchrow('SELECT * FROM users WHERE id = $1', uid)
    await c.execute('UPDATE users SET seen = now() WHERE id = $1', uid)

Fragment: assumes dsn, org and uid exist.

MethodReturns
fetchlist of Record
fetchrowone Record, or None
fetchvala single value
executea status string such as 'INSERT 0 1'

A connection runs one query at a time. If two tasks share one connection and both await it, you get an InterfaceError. Give each task its own connection with pool.fetch() (which borrows a connection just for that call) or its own pool.acquire().

Transactions and savepoints

async with c.transaction() commits on a clean exit and rolls back when an exception escapes. Pass isolation='serializable' to pick the level. A transaction block opened inside another one becomes a savepoint, so only the inner block rolls back.

text
async with c.transaction(isolation='serializable'):
    await c.execute(debit, src, amount)
    async with c.transaction():          # SAVEPOINT
        await c.execute(audit, src, amount)

Fragment: debit and audit are SQL strings.

Behind PgBouncer in transaction mode

In transaction mode PgBouncer may send the next statement to a different server connection. A prepared statement that asyncpg cached on the first connection then does not exist on the second, and you see prepared statement "__asyncpg_stmt_1__" does not exist. Turn the statement cache off. Through SQLAlchemy, two keys are needed.

text
# plain asyncpg
conn = await asyncpg.connect(dsn, statement_cache_size=0)

# through SQLAlchemy: both keys
engine = create_async_engine(
    url,
    connect_args={'statement_cache_size': 0, 'prepared_statement_cache_size': 0},
)

Fragment: url is a postgresql+asyncpg URL.

SQLAlchemy engine and session

Build the engine from a postgresql+asyncpg:// URL and wrap it in an async_sessionmaker with expire_on_commit=False, so loaded attributes stay readable after a commit instead of raising MissingGreenlet. Create both once at startup. For each request or task, open a new session with async with Session() as s, s.begin():, which commits on success and rolls back on error. Never share a session between concurrent tasks.

text
engine = create_async_engine('postgresql+asyncpg://user:pw@host/db', pool_pre_ping=True)
Session = async_sessionmaker(engine, expire_on_commit=False)

async with Session() as s, s.begin():
    res = await s.execute(select(User).where(User.age > 30))
    users = res.scalars().all()

Fragment: User is a mapped class.

Common mistake

Sharing one AsyncSession across tasks started with asyncio.gather. Open one session per task instead, and share only the engine and the sessionmaker.

Loading, Pools and Isolation

Relationships: load on purpose

An AsyncSession cannot run a hidden query when you read user.posts, because attribute access cannot be awaited. Decide how each relationship loads when you write the query. Set lazy='raise' on relationships so a surprise load fails loudly in tests instead of silently costing a query per row.

CaseUseEffect
One-to-many collectionselectinloadTwo queries, the second with IN (...)
Many-to-onejoinedloadOne query with a LEFT JOIN
Everything else by defaultlazy='raise'Fails loudly instead of loading
A deliberate load laterawaitable_attrsAwaited access to one attribute
text
stmt = select(User).options(selectinload(User.posts))
users = (await s.scalars(stmt)).all()

posts = await u.awaitable_attrs.posts   # needs the AsyncAttrs mixin

Fragment: eager-load before leaving the session.

MissingGreenlet

This error means code touched an implicit lazy load or an expired attribute. Eager-load the data in the query, or await u.awaitable_attrs.posts. If you use joinedload on a collection, call .unique() on the result.

Pool math

Every process holds its own pool, and a pool can grow past pool_size by max_overflow. The total across your deployment must stay under the server's limit, with headroom left for migrations and psql:

The budget rule

instances × workers × (pool_size + max_overflow) < max_connections. For example, 2 × 4 × (5 + 10) = 120, which is more than 100, so shrink the pool or the overflow. Also set pool_pre_ping=True so a connection killed by a restart or firewall is replaced at checkout instead of failing a request.

Isolation levels

LevelWhat a query seesWhat to do
READ COMMITTED (default)Data committed before each statement beganFine for most work
REPEATABLE READOne snapshot for the whole transactionUse for consistent reports
SERIALIZABLESnapshot plus conflict detectionRetry the whole transaction on error 40001

SQLite and Redis

SQLite allows one writer at a time, so with aiosqlite run PRAGMA journal_mode=WAL and PRAGMA busy_timeout=5000 on each connection, so waiting writers wait instead of failing with 'database is locked'. For Redis, make one shared connection pool per app, batch commands with a pipeline, and close the client at shutdown with await r.aclose().

text
import redis.asyncio as redis

pool = redis.ConnectionPool.from_url(url, max_connections=50)
r = redis.Redis(connection_pool=pool)

async with r.pipeline(transaction=True) as p:
    p.incr('hits')
    p.expire('hits', 60)
    await p.execute()                  # one round trip

await r.aclose()                       # at shutdown
await pool.disconnect()

Fragment: url is a redis:// URL. A pool you pass in is not closed for you.

Why must the pipeline be created from the shared client and not from a new client per request?
  • A new client per request opens its own connections, so a burst of traffic turns into a storm of new sockets.
  • A shared pool lets every request borrow from the same limited set of connections and return them quickly.

Timeouts, Cancellation and Retries

Put a deadline on both sides

A client-side asyncio.timeout stops your code from waiting, and asyncpg then sends a cancel request so the server stops working too. A server-side statement_timeout protects you when the client is gone or buggy. Use both, and set the server value slightly lower than the client one so the server's clear error arrives first.

text
async with asyncio.timeout(5):                      # client layer
    rows = await conn.fetch(q)

engine = create_async_engine(
    url,
    connect_args={'server_settings': {'statement_timeout': '4s', 'lock_timeout': '2s'}},
)                                                   # server layer

Fragment: server settings apply to every new pooled connection.

Keep transactions free of outside I/O

While a transaction is open it holds row locks and one pooled connection. If you await an HTTP call inside it, every other task that needs that row or connection waits for a network you do not control. Make the call first, then open the transaction and write the results. A leaked transaction shows up as idle in transaction in pg_stat_activity, and idle_in_transaction_session_timeout makes the server close it.

Never swallow cancellation

When a task is cancelled or a timeout fires, CancelledError is raised at the current await. It is a BaseException, so except Exception lets it pass, but except BaseException: pass buries it and the task refuses to stop. Do your cleanup, then re-raise. The example below runs without any database: a fake transaction fails twice with a serialization error, the whole thing is retried with backoff, and then a slow call is cut off by a deadline.

python
import asyncio

class SerializationError(Exception):
    pass

attempts = 0

async def transfer():
    global attempts
    attempts += 1
    await asyncio.sleep(0.01)
    if attempts < 3:
        raise SerializationError('40001')
    return 'committed'

async def with_retry():
    for n in range(3):
        try:
            return await transfer()
        except SerializationError:
            print(f'attempt {n + 1}: 40001, retrying')
            await asyncio.sleep(0.01 * 2**n)
    raise RuntimeError('gave up')

async def slow_query():
    try:
        await asyncio.sleep(10)
    except asyncio.CancelledError:
        print('cleanup, then re-raise')
        raise

async def main():
    print(await with_retry())
    try:
        async with asyncio.timeout(0.05):
            await slow_query()
    except TimeoutError:
        print('client deadline hit')

asyncio.run(main())
output
attempt 1: 40001, retrying
attempt 2: 40001, retrying
committed
cleanup, then re-raise
client deadline hit

Notice that the retry wraps the whole unit of work, not a single statement, because a failed serializable transaction has already been rolled back. Retry only work that is safe to repeat: a serialization failure or a lost connection qualifies, but an INSERT that timed out may have committed, so use an idempotency key before repeating it.

Common mistakes

Setting only a client timeout and no statement_timeout; awaiting HTTP inside a transaction; writing except BaseException: pass around database calls; and cleaning up without re-raising the cancellation.

When is asyncio.shield appropriate in database code?
  • Use it only for short cleanup that must finish even when the caller is cancelled, such as a final commit or releasing a connection.
  • The inner task keeps running, but the outer await still raises CancelledError, so shield protects the cleanup and never the caller.

Part 12 · Check yourself

Quiz

Each question shows a short situation. Decide what happens before you open the answer.

This handler runs under load and sometimes fails. What is the bug, and how do you fix it?
  • Both fetch calls are awaited at the same time on one connection, and a connection runs one query at a time.
  • asyncpg raises InterfaceError: another operation is in progress.
  • Give each task its own connection: use pool.fetch() for each query, or pool.acquire() once per task.
  • Sharing a connection only works when the awaits run one after another.
async def dashboard(pool):
    async with pool.acquire() as conn:
        users, orders = await asyncio.gather(
            conn.fetch('SELECT * FROM users'),
            conn.fetch('SELECT * FROM orders'),
        )
    return users, orders
The commit succeeds, but the last line crashes. What is raised, and what are two ways to prevent it?
  • MissingGreenlet. The default expire_on_commit=True expires the object at commit, so reading u.name needs a refresh. That refresh is implicit I/O, and it cannot await.
  • Fix 1: create the factory with async_sessionmaker(engine, expire_on_commit=False).
  • Fix 2: load the value deliberately, for example await s.refresh(u), before you read it.
  • The same error appears for an implicit lazy load of a relationship. Use selectinload, awaitable_attrs or lazy='raise' so the problem shows up in tests.
Session = async_sessionmaker(engine)  # default settings

async with Session() as s:
    u = User(name='ada')
    s.add(u)
    await s.commit()
    print(u.name)
A service runs 3 instances with 4 workers each. Every worker uses the SQLAlchemy defaults (pool_size=5, max_overflow=10). Postgres has max_connections=100. What happens, and what could you change?
  • Worst case is 3 x 4 x (5 + 10) = 180 connections, well above 100 minus the 3 reserved superuser slots and any admin headroom.
  • Under a traffic burst, new connections are refused with too many clients errors, and psql and migrations get locked out too.
  • Cut the per-process budget. For example pool_size=3, max_overflow=2 gives 3 x 4 x 5 = 60.
  • Or run PgBouncer and use NullPool in the app so the pooler owns pooling.
  • Check the need first with Little's law: 500 rps x 0.02 s held is only 10 connections in total.
The timeout fires after 5 seconds, but the query keeps the connection busy and the caller never sees the timeout. Why, and what is the correct handler?
  • asyncio.timeout cancels the task by raising CancelledError at the current await. CancelledError is a BaseException, and except BaseException: pass swallows it.
  • Because the cancel is swallowed, the timeout can never turn into TimeoutError. The task carries on as if nothing happened, perhaps on a connection that was cancelled mid-read.
  • Clean up, then re-raise: catch asyncio.CancelledError, do the short cleanup, and use a bare raise.
  • Pair the client deadline with a server-side statement_timeout so Postgres also stops the query.
async with asyncio.timeout(5):
    try:
        rows = await conn.fetch(slow_query)
    except BaseException:
        pass
Users report that the pool is exhausted and pg_stat_activity shows many rows in state idle in transaction, yet the database CPU is nearly idle. What is the likely bug in this code?
  • The await http.get(...) runs inside the transaction. While the HTTP call waits, the transaction stays open, its row locks are held, and one pooled connection is pinned.
  • With enough requests every connection is pinned, so other tasks queue on the pool while Postgres does nothing.
  • Move the HTTP call before conn.transaction() or after it commits, and keep the transaction short.
  • As a safety net, set idle_in_transaction_session_timeout so the server kills leaked transactions.
async with conn.transaction():
    await conn.execute(update_order, order_id)
    receipt = await http.get(PAYMENT_URL)
    await conn.execute(insert_receipt, order_id, receipt.text)

Summary

  • Async drivers (asyncpg, aiosqlite, redis.asyncio) let the loop serve other tasks while the database works. They add concurrency, not speed, and a blocking call in async def freezes everything.
  • One asyncpg connection runs one query at a time. Share a pool, not a connection, and use $1 placeholders so values never touch the SQL text.
  • Behind PgBouncer in transaction mode, turn the statement caches off (statement_cache_size=0, plus prepared_statement_cache_size=0 through SQLAlchemy).
  • Use one AsyncSession per request or task, with expire_on_commit=False. Eager-load with selectinload or joinedload, and treat MissingGreenlet as a missed load.
  • Keep transactions short and free of external I/O. Retry the whole transaction on SerializationFailure (40001).
  • Size pools with instances x workers x (pool_size + max_overflow) below max_connections, and use Little's law instead of guessing big.
  • Set deadlines on both sides (asyncio.timeout and statement_timeout), never swallow CancelledError, and retry only idempotent work.