Python 基础体系 · 第 85/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 异步数据库访问:连接池、事务、取消、并发和一致性

异步数据库访问的难点不在于把 def 改成 async def,而在于同时处理五个相互影响的问题:

  1. 连接池:有限的数据库连接如何被大量请求复用。
  2. 事务:多条 SQL 如何构成一个不可分割的业务操作。
  3. 取消:客户端断开、请求超时或服务关闭时,正在执行的数据库操作如何收尾。
  4. 并发:多个协程、多个请求和多个进程如何同时访问数据库。
  5. 一致性:在并发读写、重试和部分失败下,系统如何避免错误数据。

这五个问题不能分别孤立处理。例如,请求取消可能发生在事务提交前;连接池耗尽可能表现为接口超时;把同一个 AsyncSession 交给多个并发任务,可能破坏会话的状态机;一个看似正确的“先查询、再更新”流程,也可能在并发下丢失更新。

本文以 Python 3.14 的 asyncio 和 SQLAlchemy 2.0 异步 API 为背景,使用 FastAPI 风格的 ASGI 应用说明完整生命周期。Python 3.14 的 asyncio 仍然采用协作式调度:事件循环一次运行一个任务,任务在等待 Future 或 I/O 时让出执行权,其他任务才有机会运行。(docs.python.org)


一、先建立正确的模型:异步不等于数据库并行

1.1 协程、任务和数据库 I/O

协程函数只是一个可等待对象:

async def load_user(user_id: int):
    return await session.get(User, user_id)

调用 load_user(1) 并不会立即执行函数体,而是得到一个 coroutine object。只有以下操作之一才会推动它执行:

await load_user(1)

或者:

task = asyncio.create_task(load_user(1))
await task

Task 是由事件循环调度的协程执行单元。异步数据库驱动在等待网络响应时会把控制权交还给事件循环,因此同一个线程可以在等待数据库期间处理其他请求。

但是,下面这段代码并不会因为使用了 asyncio 就自动变成数据库并行:

async with session.begin():
    user = await session.get(User, 1)
    user.name = "new-name"
    await session.flush()
    result = await session.execute(select(Order))

同一个事务中的数据库命令仍然需要按照顺序发送到数据库连接。SQLAlchemy 文档将 SessionAsyncSession 和数据库连接都描述为有状态对象;一个会话代表一个逻辑事务,并且同一 AsyncSession 不适合被多个 asyncio 任务同时使用。(docs.sqlalchemy.org)

因此应区分三种“并发”:

层次 含义 是否能并行
协程并发 多个任务交错执行 可以
数据库连接并发 多个连接同时向数据库发请求 可以
同一连接上的命令 在同一事务连接上交错操作 通常不应并行

异步 API 解决的是“等待数据库时不阻塞事件循环”,不是“让一个数据库连接同时执行多条 SQL”。


二、连接池:控制数据库连接的并发入口

2.1 连接池解决什么问题

数据库连接不是普通的 Python 对象。建立连接通常需要:

  1. 创建 TCP 连接;
  2. 完成数据库协议握手;
  3. 认证;
  4. 设置会话参数;
  5. 可能执行初始化 SQL。

如果每个 HTTP 请求都新建连接,连接建立成本会重复发生,更重要的是,突发流量会直接把数据库连接数推高。

连接池把连接生命周期改成:

创建少量连接
      ↓
请求到来
      ↓
从池中借出连接
      ↓
执行事务和 SQL
      ↓
归还连接

连接池的容量不是“应用能处理的请求数”,而是“同一进程可以同时占用的数据库连接数”。

设:

  • W:应用进程或 worker 数量;
  • P:每个进程的池大小;
  • O:每个进程允许额外创建的溢出连接;
  • C_db:数据库允许该应用使用的最大连接数。

应用理论上最多可能占用:

Capp=W×(P+O)C_{app} = W \times (P + O)

要避免应用自身突破数据库连接上限,至少需要满足:

W×(P+O)+CotherCdbW \times (P + O) + C_{other} \leq C_{db}

其中 C_other 表示管理工具、迁移程序、监控组件和其他服务占用的连接。

例如:

  • 4 个 worker;
  • 每个 worker pool_size=10
  • 每个 worker max_overflow=5
  • 数据库为该服务预留 70 个连接;
  • 其他组件需要 5 个连接。

则:

4×(10+5)+5=654 \times (10 + 5) + 5 = 65

在这个简单估算下还有 5 个连接余量。但这只是上限估算,不代表所有请求都能立即获得连接。

2.2 池大小和请求并发的关系

假设连接池只有 10 个连接,同时有 100 个请求执行数据库操作:

请求 1  ─┐
请求 2  ─┤
...
请求 10 ─┤── 获得连接,执行 SQL
请求 11 ─┤
...
请求 100─┘── 等待连接归还

第 11 个请求并不是自动失败,而是通常等待池中出现可用连接。等待时间过长时,最终表现为:

  • 请求超时;
  • 连接池超时异常;
  • 上游重试;
  • 更多请求进入等待;
  • 数据库和应用同时出现级联压力。

连接池因此是一个并发闸门。它把无限请求并发限制为有限数据库并发。

但池越大不一定越快。数据库执行查询、锁竞争、磁盘 I/O 和 CPU 资源都有限。过大的池可能导致更多连接同时争抢数据库资源,反而增加上下文切换、锁等待和内存消耗。

2.3 Engine、连接池和 AsyncSession 的职责

在 SQLAlchemy 异步架构中,可以这样划分职责:

AsyncEngine
 ├── 管理连接池
 ├── 创建 AsyncConnection
 └── 连接数据库驱动

AsyncSession
 ├── 管理 ORM 对象状态
 ├── 组织事务
 ├── 发出 ORM/Core SQL
 └── 使用 AsyncEngine 获取连接

AsyncEngine 通常是应用级对象,一个进程创建一个;AsyncSession 通常是请求级对象,一个请求或一个独立业务单元创建一个。

不要把 AsyncSession 做成全局单例:

# 错误示例
global_session = AsyncSessionFactory()

async def handler():
    return await global_session.execute(...)

原因不是“全局变量风格不好”,而是 AsyncSession 内部包含当前事务、对象身份映射、待刷新的对象和连接状态。多个请求同时修改它们,会让不同业务操作共享一个本不该共享的状态机。


三、一个可运行的异步数据库应用

下面使用 SQLite 和 aiosqlite,便于本地运行。生产环境通常会替换为 PostgreSQL 或 MySQL 对应的异步驱动,但连接池、事务和会话隔离的基本原则不变。

安装依赖:

python -m pip install "fastapi" "uvicorn[standard]" "sqlalchemy>=2.0" "aiosqlite"

目录结构:

app.py

完整示例:

from __future__ import annotations

from contextlib import asynccontextmanager
from typing import AsyncIterator

from fastapi import Depends, FastAPI, HTTPException
from pydantic import BaseModel, Field
from sqlalchemy import String, select
from sqlalchemy.ext.asyncio import (
    AsyncSession,
    async_sessionmaker,
    create_async_engine,
)
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column


DATABASE_URL = "sqlite+aiosqlite:///./demo.db"


class Base(DeclarativeBase):
    pass


class Account(Base):
    __tablename__ = "account"

    id: Mapped[int] = mapped_column(primary_key=True)
    owner: Mapped[str] = mapped_column(String(100), nullable=False)
    balance: Mapped[int] = mapped_column(nullable=False, default=0)


class TransferRequest(BaseModel):
    source_id: int
    target_id: int
    amount: int = Field(gt=0)


engine = create_async_engine(
    DATABASE_URL,
    echo=True,
    pool_pre_ping=True,
)

SessionFactory = async_sessionmaker(
    bind=engine,
    class_=AsyncSession,
    expire_on_commit=False,
)


@asynccontextmanager
async def lifespan(app: FastAPI):
    async with engine.begin() as connection:
        await connection.run_sync(Base.metadata.create_all)

    yield

    await engine.dispose()


app = FastAPI(lifespan=lifespan)


async def get_session() -> AsyncIterator[AsyncSession]:
    async with SessionFactory() as session:
        yield session


@app.post("/accounts")
async def create_account(
    owner: str,
    balance: int = 0,
    session: AsyncSession = Depends(get_session),
):
    account = Account(owner=owner, balance=balance)
    session.add(account)

    await session.commit()
    await session.refresh(account)

    return {
        "id": account.id,
        "owner": account.owner,
        "balance": account.balance,
    }


@app.get("/accounts/{account_id}")
async def get_account(
    account_id: int,
    session: AsyncSession = Depends(get_session),
):
    account = await session.get(Account, account_id)

    if account is None:
        raise HTTPException(status_code=404, detail="account not found")

    return {
        "id": account.id,
        "owner": account.owner,
        "balance": account.balance,
    }


@app.post("/transfers")
async def transfer(
    request: TransferRequest,
    session: AsyncSession = Depends(get_session),
):
    if request.source_id == request.target_id:
        raise HTTPException(status_code=400, detail="accounts must differ")

    async with session.begin():
        source = await session.get(Account, request.source_id)
        target = await session.get(Account, request.target_id)

        if source is None or target is None:
            raise HTTPException(status_code=404, detail="account not found")

        if source.balance < request.amount:
            raise HTTPException(status_code=409, detail="insufficient balance")

        source.balance -= request.amount
        target.balance += request.amount

    return {"status": "committed"}

启动:

uvicorn app:app --reload

创建账户:

curl -X POST "http://127.0.0.1:8000/accounts?owner=Alice&balance=100"
curl -X POST "http://127.0.0.1:8000/accounts?owner=Bob&balance=0"

转账:

curl -X POST \
  -H "Content-Type: application/json" \
  -d '{"source_id": 1, "target_id": 2, "amount": 30}' \
  "http://127.0.0.1:8000/transfers"

预期结果:

{"status":"committed"}

再次查询时,Alice 的余额为 70,Bob 的余额为 30。

这里有三个生命周期:

应用启动
  └── 创建 AsyncEngine
  └── 建表或执行初始化逻辑

请求开始
  └── 创建 AsyncSession
  └── 执行查询或事务
  └── 请求结束时关闭 AsyncSession

应用关闭
  └── 等待生命周期退出
  └── dispose AsyncEngine

FastAPI 的 lifespan 适合管理应用级资源;请求依赖适合管理请求级会话。ASGI 规范把应用描述为由协议服务器调用的异步应用,并通过事件消息传递请求、断开和响应状态;HTTP 连接断开可以通过 http.disconnect 事件通知应用,但在高并发场景下,发送响应时也可能先收到服务器特定的 OSError。(asgi.readthedocs.io)


四、事务:把业务不变量绑定到提交边界

4.1 什么是事务

事务是一组数据库操作的逻辑单元。常见的四个属性是 ACID:

  • Atomicity,原子性:全部成功,或全部回滚;
  • Consistency,一致性:事务前后满足数据库约束和业务不变量;
  • Isolation,隔离性:并发事务之间按照数据库隔离规则观察彼此;
  • Durability,持久性:提交成功后,数据不会因为普通进程故障而丢失。

事务并不自动保证业务正确。数据库只能保证声明出来的约束和所选隔离级别。比如“账户余额不能为负数”如果只写在 Python 判断里,而没有锁、条件更新或数据库约束保护,就可能在并发下失效。

4.2 转账的原子性

转账需要满足不变量:

balancesource=balancesourceamountbalance_{source}' = balance_{source} - amount

balancetarget=balancetarget+amountbalance_{target}' = balance_{target} + amount

两条更新必须同时发生,否则可能出现:

扣款成功
进程崩溃
入账没有执行

事务将两条更新绑定到同一个提交点:

async with session.begin():
    source.balance -= amount
    target.balance += amount

退出 async with 时:

  • 没有异常:提交事务;
  • 发生异常:回滚事务;
  • 事务提交失败:向调用者抛出异常,调用者不能把它当成成功。

因此不能这样写:

async with session.begin():
    source.balance -= amount

await external_payment_service()
target.balance += amount
await session.commit()

这段代码先提交了扣款,再调用外部服务,最终无法让数据库事务覆盖整个跨系统操作。更合理的设计通常是:

  1. 在一个数据库事务中创建支付意图或待处理记录;
  2. 提交;
  3. 调用外部服务;
  4. 通过幂等键和状态机记录结果;
  5. 失败时由补偿任务重试。

数据库事务不能跨越任意外部系统获得原子性。

4.3 commit()flush() 的区别

flush() 把当前会话中的变更发送给数据库,但不结束事务:

account.balance -= 10
await session.flush()

这时数据库可能已经执行了 UPDATE,但其他事务通常仍然看不到最终提交结果,具体可见性取决于数据库和隔离级别。

commit() 则执行提交,使事务中的变更成为持久状态:

await session.commit()

典型流程是:

Python 对象变化
  ↓
flush
  ↓
数据库执行 INSERT / UPDATE / DELETE
  ↓
commit
  ↓
事务持久化

async with session.begin() 中,SQLAlchemy 会在退出上下文时负责提交或回滚。不要在同一个上下文中随意再次 commit(),否则会破坏事务边界的可读性。

4.4 异常路径

事务代码必须显式考虑异常:

async def transfer(
    session: AsyncSession,
    source_id: int,
    target_id: int,
    amount: int,
) -> None:
    async with session.begin():
        source = await session.get(Account, source_id)
        target = await session.get(Account, target_id)

        if source is None or target is None:
            raise LookupError("account not found")

        if source.balance < amount:
            raise ValueError("insufficient balance")

        source.balance -= amount
        target.balance += amount

如果 source.balance -= amount 执行后,第二次赋值前发生异常,事务退出时仍然会回滚已经 flush 的修改。

反例:

async def unsafe_transfer(session, source_id, target_id, amount):
    source = await session.get(Account, source_id)
    source.balance -= amount
    await session.commit()

    target = await session.get(Account, target_id)
    target.balance += amount
    await session.commit()

这里有两个独立事务。第二次提交失败时,第一笔扣款已经不可回滚。


五、并发:不要共享 AsyncSession,要共享 Engine

5.1 错误的共享方式

async def load_two_accounts(session: AsyncSession):
    async def load(account_id: int):
        return await session.get(Account, account_id)

    return await asyncio.gather(
        load(1),
        load(2),
    )

问题不在 gather() 本身,而在两个任务共享同一个 AsyncSession。会话包含事务和 ORM 状态,两个任务同时发出命令会让操作顺序、事务状态和对象状态变得难以推断。

SQLAlchemy 明确要求:一个 AsyncSession 不应被多个 asyncio 任务同时使用;并发任务应该使用各自的 AsyncSession。(docs.sqlalchemy.org)

正确方式是每个任务创建自己的会话:

async def load_account(account_id: int) -> dict:
    async with SessionFactory() as session:
        account = await session.get(Account, account_id)

        if account is None:
            raise LookupError(account_id)

        return {
            "id": account.id,
            "owner": account.owner,
            "balance": account.balance,
        }


async def load_accounts() -> list[dict]:
    async with asyncio.TaskGroup() as group:
        task_a = group.create_task(load_account(1))
        task_b = group.create_task(load_account(2))

    return [task_a.result(), task_b.result()]

这里有:

Task A ── Session A ── Connection A
Task B ── Session B ── Connection B

是否真的同时执行,还受到连接池大小和数据库服务能力限制。如果连接池只有一个可用连接,两个任务仍然会在连接获取阶段排队。

5.2 gather()TaskGroup

asyncio.gather() 可以并发等待多个可等待对象,但默认异常语义容易被误解:一个任务失败时,其他任务不一定自动取消。

TaskGroup 是结构化并发 API。任务组退出时会等待其子任务;当某个子任务抛出非 CancelledError 异常时,其他任务会被取消,并将异常集中处理。Python 文档将 TaskGroup 作为比 gather() 更安全的任务组织方式之一。(docs.python.org)

读取多个互不相关的数据源时,可以使用:

async def dashboard() -> dict:
    async with asyncio.TaskGroup() as group:
        user_task = group.create_task(load_user())
        orders_task = group.create_task(load_orders())
        notices_task = group.create_task(load_notices())

    return {
        "user": user_task.result(),
        "orders": orders_task.result(),
        "notices": notices_task.result(),
    }

但不要为了“并发”把同一个数据库事务拆成多个会话:

事务 A:扣款
事务 B:入账

这会失去转账所需的原子性。并行读取和并行写入的约束不同:

  • 独立读取:可以使用不同会话并发执行;
  • 同一业务事务:必须保持一个明确事务边界;
  • 互相依赖的写入:通常应在一个事务中顺序执行;
  • 大量独立写入:应评估批量 SQL、数据库锁和连接池容量。

5.3 N+1 查询在异步环境中仍然存在

异步不会消除 N+1 查询。

users = (await session.execute(select(User))).scalars().all()

for user in users:
    print(user.orders)

如果 orders 是延迟加载关系,循环中的每次访问都可能触发一次查询:

1 次查询:获取用户
N 次查询:分别获取每个用户的订单
总计 N + 1 次

异步 ORM 还会让隐式 I/O 更危险,因为普通属性访问不是 await 点。应使用显式加载策略,例如:

from sqlalchemy import select
from sqlalchemy.orm import selectinload

stmt = (
    select(User)
    .options(selectinload(User.orders))
)

result = await session.execute(stmt)
users = result.scalars().all()

这通常会把查询组织为:

查询用户
查询这些用户对应的全部订单

而不是每个用户单独查询。SQLAlchemy 官方异步 ORM 示例也专门区分了异步 ORM、批量查询和关系加载方式。(docs.sqlalchemy.org)


六、取消:请求结束不代表数据库操作必然停止

6.1 CancelledError 的传播

取消不是普通返回值。调用:

task.cancel()

会请求取消任务;任务在下一个可取消的等待点收到 asyncio.CancelledError。这个异常直接继承自 BaseException,不是普通 Exception 的子类,因此不应使用宽泛的异常处理把它吞掉。(docs.python.org)

数据库代码必须使用 finally 做资源清理:

async def query_data(session: AsyncSession):
    try:
        result = await session.execute(...)
        return result
    finally:
        # 这里执行必要的本地清理
        pass

请求依赖中的:

async def get_session():
    async with SessionFactory() as session:
        yield session

在请求结束或任务取消时退出上下文,能够释放会话并归还连接。

错误写法:

try:
    result = await session.execute(stmt)
except BaseException:
    return []

这会把取消、系统退出等控制流异常伪装成正常结果,可能造成:

  • 事务没有按预期回滚;
  • 连接没有及时归还;
  • 上层无法知道请求已取消;
  • 结构化并发失效。

如果捕获 CancelledError 是为了记录日志或执行清理,应在清理完成后重新抛出:

import asyncio
import logging

logger = logging.getLogger(__name__)


async def work(session: AsyncSession):
    try:
        return await session.execute(...)
    except asyncio.CancelledError:
        logger.info("database work cancelled")
        raise

Python 3.14 文档特别提醒,TaskGroupasyncio.timeout() 内部依赖取消机制;吞掉 CancelledError 可能导致它们行为异常。(docs.python.org)

6.2 超时和取消的区别

超时是业务策略,取消是执行控制流。

import asyncio


async def load_with_timeout(session: AsyncSession):
    try:
        async with asyncio.timeout(2):
            return await session.execute(...)
    except TimeoutError:
        raise RuntimeError("database operation timed out")

这里 asyncio.timeout(2) 会在 2 秒后取消当前代码块,随后把内部取消转换为 TimeoutError。因此可以在外层按业务含义处理超时。

但超时只保证 Python 任务不再继续等待,不保证数据库服务器已经停止执行 SQL。可能出现以下时间线:

T0  客户端发送请求
T1  应用向数据库发送 UPDATE
T2  应用侧超时,任务被取消
T3  数据库仍在执行 UPDATE
T4  数据库完成并提交,或因连接关闭回滚

是否回滚取决于数据库协议、驱动行为、连接是否处于事务中以及取消发生的具体位置。应用不能仅凭“Python 抛出了 TimeoutError”推断数据库端一定没有发生写入。

对写操作尤其要使用幂等设计。例如创建订单时携带业务唯一键:

CREATE UNIQUE INDEX uq_order_request
ON orders(request_id);

重试时再次使用同一个 request_id,数据库会阻止重复创建,应用再查询已有结果即可。

6.3 客户端断开和数据库事务

ASGI HTTP 规范定义了 http.disconnect 事件,用于通知应用连接已断开;如果应用继续向关闭的连接发送数据,服务器可能抛出 OSError。规范还提醒,在高并发代码中,发送失败可能早于收到 http.disconnect。(asgi.readthedocs.io)

这意味着:

客户端断开
   ↓
HTTP 响应没有意义
   ↓
数据库工作是否继续,需要应用自己决定

不要把“客户端已经断开”简单等同于“数据库必须取消”。有些操作应该尽快停止,例如大查询、报表生成和非必要预加载;有些操作不能因为客户端断开就撤销,例如已经接受的支付、库存扣减或审计记录。

核心判断是:这项数据库工作属于哪一类?

  • 响应型工作:只为当前请求服务,客户端断开后通常可以取消;
  • 状态型工作:代表已经接受的业务命令,即使客户端断开也应继续完成;
  • 补偿型工作:由后台任务、消息队列或定时任务继续执行。

如果业务命令需要可靠完成,不应依赖 HTTP 请求任务本身。应将命令写入数据库或消息系统后,再由独立 worker 执行。


七、一致性:从业务不变量推导并发控制

7.1 “先查再改”的竞态

考虑两个并发请求同时从账户扣款:

account = await session.get(Account, 1)

if account.balance >= 80:
    account.balance -= 80
    await session.commit()

初始余额为 100。

两个事务可能按如下顺序运行:

事务 A:读取余额 100
事务 B:读取余额 100
事务 A:计算 20
事务 B:计算 20
事务 A:写入 20
事务 B:写入 20

最终余额是 20,但实际扣款总额是 160。系统丢失了一次更新,也违反了余额不应透支的业务不变量。

这说明:

读取满足条件
+
本地计算
+
稍后写回

不是一个原子操作。

7.2 条件更新:让数据库一次完成判断和修改

可以把判断放进 UPDATE

UPDATE account
SET balance = balance - :amount
WHERE id = :account_id
  AND balance >= :amount;

然后检查影响行数:

from sqlalchemy import update


async def withdraw(
    session: AsyncSession,
    account_id: int,
    amount: int,
) -> bool:
    async with session.begin():
        stmt = (
            update(Account)
            .where(
                Account.id == account_id,
                Account.balance >= amount,
            )
            .values(balance=Account.balance - amount)
        )

        result = await session.execute(stmt)

        return result.rowcount == 1

这里的逻辑是:

  • rowcount == 1:条件成立且扣款成功;
  • rowcount == 0:账户不存在或余额不足;
  • 整个判断与修改由数据库作为一条语句执行。

这比把余额读取到 Python 后再判断更接近业务原子性。

7.3 悲观锁:锁住将要修改的行

另一种方式是使用行锁:

from sqlalchemy import select


async def withdraw_with_lock(
    session: AsyncSession,
    account_id: int,
    amount: int,
) -> bool:
    async with session.begin():
        stmt = (
            select(Account)
            .where(Account.id == account_id)
            .with_for_update()
        )

        account = (await session.execute(stmt)).scalar_one_or_none()

        if account is None or account.balance < amount:
            return False

        account.balance -= amount
        return True

典型执行过程:

事务 A:SELECT ... FOR UPDATE,获得锁
事务 B:SELECT ... FOR UPDATE,等待
事务 A:读取 100,扣款到 20,提交
事务 B:获得锁,读取最新余额 20
事务 B:发现余额不足,回滚或返回失败

FOR UPDATE 的语义依赖数据库和存储引擎。SQLite 对行级锁的支持与 PostgreSQL、MySQL 不同,不能把 PostgreSQL 的锁行为直接推断到 SQLite。

7.4 乐观锁:使用版本号检测冲突

乐观锁适合冲突不频繁、希望减少锁等待的场景。表中增加版本字段:

ALTER TABLE account ADD COLUMN version INTEGER NOT NULL DEFAULT 0;

更新时同时检查旧版本:

UPDATE account
SET balance = :new_balance,
    version = version + 1
WHERE id = :account_id
  AND version = :old_version;

如果影响行数为 0,说明其他事务已经修改了这条记录。应用可以:

  1. 重新读取;
  2. 重新计算;
  3. 重试有限次数;
  4. 或向调用者报告冲突。

乐观锁的核心公式是:

成功    id=iversion=v\text{成功} \iff id = i \land version = v

更新后:

version=v+1version' = v + 1

任何并发事务只要仍然使用旧版本 v,就不能再次成功更新。


八、隔离级别:一致性不是一个开关

隔离级别规定一个事务能观察到哪些并发变化。常见概念包括:

  • 脏读:读到了其他事务尚未提交的数据;
  • 不可重复读:同一事务两次读取同一行,结果不同;
  • 幻读:同一条件查询两次,第二次出现了新增或消失的行;
  • 丢失更新:两个事务基于旧值写回,后写覆盖先写。

隔离级别越强,通常越能减少异常现象,但可能增加锁等待、冲突和重试成本。具体实现由数据库决定,不能仅依赖 ORM API 的名称推断真实行为。

例如,在一个事务中:

async with session.begin():
    first = await session.scalar(
        select(Account.balance).where(Account.id == 1)
    )

    await asyncio.sleep(1)

    second = await session.scalar(
        select(Account.balance).where(Account.id == 1)
    )

firstsecond 是否相同,不由 Python 的 async 关键字决定,而由:

  1. 数据库隔离级别;
  2. 查询是否使用同一事务;
  3. 其他事务何时提交;
  4. 数据库的 MVCC 或锁实现;

共同决定。

一个常见误解是:“使用事务就不会有并发问题。”实际上事务只提供边界,具体保护效果还需要:

事务边界
+
正确的隔离级别
+
锁或条件更新
+
数据库约束
+
正确的重试策略

九、事务重试:只重试可安全重放的操作

数据库可能因死锁、序列化冲突、临时网络错误而失败。重试不是简单地把整个函数包在 for 循环里。

错误示例:

for _ in range(3):
    try:
        await transfer(...)
        break
    except Exception:
        continue

问题包括:

  • 已经提交成功但响应丢失,重试可能重复执行;
  • 捕获了编程错误;
  • 没有等待退避;
  • 没有区分可重试和不可重试异常;
  • 没有保证每次重试使用新的事务。

相对安全的结构是:

import asyncio


async def run_transaction_with_retry(
    operation,
    attempts: int = 3,
):
    for attempt in range(attempts):
        try:
            return await operation()
        except TransientDatabaseError:
            if attempt + 1 == attempts:
                raise

            await asyncio.sleep(0.05 * (2**attempt))

其中 operation() 每次都必须创建新的事务上下文:

async def one_attempt() -> None:
    async with SessionFactory() as session:
        async with session.begin():
            # 只放可重放、具有幂等语义的数据库操作
            ...

如果操作包含发送邮件、扣除第三方余额或调用外部 HTTP 服务,就不能仅靠数据库事务重试保证安全。应使用:

  • 业务幂等键;
  • 唯一约束;
  • 状态机;
  • Outbox 表;
  • 消息队列;
  • 补偿操作。

十、会话过期、懒加载和异步 ORM 的边界

10.1 expire_on_commit

示例中使用:

SessionFactory = async_sessionmaker(
    bind=engine,
    expire_on_commit=False,
)

这样提交后,已经加载到对象中的字段仍可直接访问:

await session.commit()
print(account.balance)

如果使用默认的过期行为,提交后 ORM 可能将对象标记为过期;再次访问属性时需要从数据库重新加载。在异步环境中,普通属性访问并不是显式的 await,因此隐式 I/O 容易造成错误或不可预测的查询。

这不是说 expire_on_commit=False 永远正确,而是要明确选择:

  • 数据提交后只返回当前对象:可以关闭过期;
  • 需要重新读取数据库真实状态:显式 refresh()
  • 依赖关系数据:显式使用 selectinload() 等加载方式;
  • 不希望访问属性触发 I/O:避免隐式懒加载。

10.2 事务外访问 ORM 对象

不应把仍然依赖会话的 ORM 对象直接跨层传递:

account = await session.get(Account, 1)
await session.commit()

return account.orders  # 可能触发额外数据库访问

更清晰的做法是,在事务和会话仍然有效时完成加载,然后转换为普通数据:

return {
    "id": account.id,
    "owner": account.owner,
    "orders": [
        {"id": order.id}
        for order in account.orders
    ],
}

这样响应层不会意外触发数据库 I/O,也不会把 ORM 会话生命周期扩展到序列化阶段。


十一、数据库连接泄漏:最容易被忽略的故障路径

连接泄漏通常不是“忘记关闭 Engine”,而是请求结束后某个会话、结果集或事务仍然占用连接。

典型故障过程:

请求 1 获取连接
请求 1 发生异常
异常路径没有退出 session.begin()
连接没有归还
      ↓
请求 2、请求 3 继续获取连接
      ↓
池中连接逐渐减少
      ↓
所有请求等待连接
      ↓
接口大量超时

应使用异步上下文管理器:

async with SessionFactory() as session:
    async with session.begin():
        ...

不要手动拼接复杂的 try/finally,除非确实需要特殊资源管理。

当出现连接池耗尽时,应检查:

  1. 是否存在未关闭的 AsyncSession
  2. 是否有事务打开后长时间等待外部 I/O;
  3. 是否把数据库查询结果交给了长生命周期任务;
  4. 是否在应用启动时重复创建 Engine;
  5. worker 数量乘以池大小是否超过数据库限制;
  6. 是否有慢查询或锁等待,使连接长时间无法归还;
  7. 是否发生了取消异常但清理代码没有执行。

pool_pre_ping=True 可以帮助检测连接是否仍然有效,但它不能解决事务泄漏、慢查询和连接池配置过大的问题。


十二、不要在数据库事务中等待无关的外部 I/O

看似自然的代码:

async with session.begin():
    order = Order(...)
    session.add(order)

    await payment_gateway.charge(...)

会让数据库事务在等待支付网关期间一直占用连接。假设外部服务响应缓慢,大量请求就会变成:

连接池连接
  ↓
开启事务
  ↓
等待外部 HTTP
  ↓
连接长期占用
  ↓
其他请求无法获取连接

更合理的流程取决于业务语义。

方案一:先创建待支付订单

事务 1:
  创建 order(status='pending', request_id=...)
  提交

调用支付服务

事务 2:
  更新 order(status='paid' 或 'failed')
  提交

方案二:Outbox

在同一个数据库事务中写入业务数据和待发送事件:

事务:
  写入订单
  写入 outbox_event
  提交

后台 worker:
  读取 outbox_event
  调用外部服务
  标记事件完成

这样数据库中的业务状态和“需要发送的事件”具有同一提交边界。外部调用失败时,可以继续重试,而不必让 HTTP 请求长期持有数据库连接。


十三、取消、回滚和连接归还的完整时序

一个请求型数据库操作的正常和异常路径可以表示为:

sequenceDiagram
    participant C as Client
    participant A as ASGI/FastAPI
    participant S as AsyncSession
    participant P as Pool
    participant D as Database

    C->>A: HTTP request
    A->>S: create session
    S->>P: acquire connection
    P->>D: BEGIN / SQL
    D-->>P: result
    P-->>S: result
    S->>D: COMMIT
    D-->>S: committed
    S->>P: release connection
    A-->>C: HTTP response

    Note over C,A: client disconnects or timeout
    A--xS: task cancellation
    S->>D: rollback if transaction is active
    S->>P: release or invalidate connection

关键点有三个:

  1. 取消发生在任意 await,可能发生在获取连接、执行 SQL、读取结果或提交阶段;
  2. 回滚不是业务成功的反义词,而是事务清理动作
  3. 连接归还必须发生,即使请求没有正常返回响应

取消处理代码应保持短小:

async def service(session: AsyncSession):
    try:
        async with session.begin():
            ...
    except asyncio.CancelledError:
        # 记录必要上下文,但不要把取消转换成成功
        raise

如果需要在取消后完成极少量不可中断清理,应谨慎使用 asyncio.shield()shield() 只保护被包装的等待对象不立即受到外层取消,并不能让整个请求重新变成可取消;误用可能延长关闭时间或继续占用连接。


十四、读写分离和一致性延迟

生产系统可能把读请求发送到只读副本,把写请求发送到主库:

写请求 → Primary
读请求 → Replica

这样会引入复制延迟:

T0:主库提交新订单
T1:接口返回成功
T2:副本尚未同步
T3:紧接着读取副本,查不到订单

这不是 ORM 缓存问题,而是数据流经过不同数据库节点造成的可见性差异。

需要“写后读一致”的场景可以:

  • 提交后短时间继续读主库;
  • 根据请求上下文携带读主库标记;
  • 使用复制位点或数据库提供的一致性机制;
  • 对最终一致的页面显示“处理中”状态;
  • 不把副本查询结果当作刚提交数据的立即证明。

读写分离改变了系统的一致性模型,必须体现在接口语义中。不能一边承诺“创建成功后立即可查询”,一边无条件把后续查询路由到可能滞后的副本。


十五、诊断:从现象反推机制

15.1 连接池耗尽

常见日志表现:

QueuePool limit ... reached
TimeoutError: QueuePool limit ...

诊断顺序:

  1. 统计请求进入数据库层的数量;
  2. 查看连接池已签出连接数;
  3. 检查事务持续时间;
  4. 查看数据库锁等待和慢查询;
  5. 检查取消路径是否释放会话;
  6. 按 worker 数量重新计算最大连接数。

如果数据库端连接数很低,但应用请求全部超时,问题可能是应用进程自己的池配置;如果数据库端连接数已经达到上限,问题可能是 worker、池大小或其他服务共同造成的。

15.2 查询变慢但 CPU 不高

异步服务中,CPU 不高不代表系统健康。大量任务可能都在等待:

await 获取连接
await 锁
await 数据库响应
await 结果集

应区分:

  • 获取连接耗时;
  • SQL 执行耗时;
  • 锁等待耗时;
  • 结果传输耗时;
  • Python 序列化耗时。

只记录整个请求耗时,无法判断到底是池不足、数据库慢还是响应序列化慢。

15.3 数据偶发错误

并发一致性问题往往具有“偶发”特征,例如:

  • 余额偶尔变负;
  • 库存偶尔超卖;
  • 重复创建订单;
  • 查询到刚刚提交的数据为空;
  • 重试后出现重复扣款。

诊断需要记录:

  • 业务请求 ID;
  • 事务开始和提交时间;
  • 数据库连接标识;
  • 影响行数;
  • 版本号或锁等待;
  • 是否发生取消和重试;
  • 数据库节点是主库还是副本。

没有这些信息,只看最终异常通常无法重建并发顺序。


十六、常见误解与正确边界

误解一:async def 会让同步数据库驱动变成异步

不会。同步驱动在 await 外部执行时仍可能阻塞事件循环。必须使用真正支持异步 I/O 的数据库驱动,或者显式把同步数据库调用放入线程池,并接受线程调度和连接管理的额外成本。

误解二:连接池越大吞吐越高

不一定。连接池只是允许更多并发数据库操作,数据库本身的 CPU、锁、磁盘和缓存仍然有限。池过大可能把排队位置从应用转移到数据库,造成更难诊断的锁等待。

误解三:每个任务都可以共享一个 Session

不可以。共享的应该是 AsyncEngine 或会话工厂;独立并发任务应该拥有独立的 AsyncSession。(docs.sqlalchemy.org)

误解四:客户端断开后,数据库写入一定被取消

不一定。HTTP 任务取消和数据库服务器执行状态不是同一个状态机。对于必须完成的命令,应使用后台任务、消息队列或持久化状态机。

误解五:事务自动解决丢失更新

不解决。事务需要结合条件更新、行锁、版本号、数据库约束或合适的隔离级别,才能保护具体业务不变量。

误解六:捕获所有异常再返回默认值更稳定

这会隐藏取消、连接错误、事务失败和编程错误。尤其不要捕获 BaseException 后返回正常响应。取消应清理后继续传播,业务异常应转换为明确的 HTTP 错误,基础设施异常应记录并交给统一错误处理。


十七、一个可执行的设计检查过程

设计一个异步数据库接口时,可以按下面的因果顺序检查,而不是先决定“用多少个连接”。

第一步:写出业务不变量

例如转账:

source.balance 不得小于 0
source 和 target 的总余额不变
同一个 request_id 不能重复处理

第二步:确定事务边界

读取账户
判断余额
扣款
入账
写入操作记录

这些操作是否必须全部成功?如果是,就放在同一事务中。

第三步:确定并发控制

  • 单条条件更新;
  • SELECT ... FOR UPDATE
  • 乐观锁版本号;
  • 唯一约束;
  • 业务幂等键。

第四步:确定会话归属

一个请求 / 一个独立事务单元 → 一个 AsyncSession
多个并发任务 → 多个 AsyncSession
整个应用 → 一个 AsyncEngine

第五步:确定取消语义

问清楚:

客户端断开后,这个操作还能不能停止?

如果不能停止,就不要把它绑定在请求任务上。

第六步:确定重试语义

问清楚:

如果提交结果未知,再执行一次会发生什么?

如果可能重复扣款、重复发货或重复创建,就必须先加入幂等键和状态查询。

第七步:计算连接上限

使用:

W×(P+O)+CotherCdbW \times (P + O) + C_{other} \leq C_{db}

然后再结合平均事务时间、慢查询、锁等待和峰值并发进行压测。公式只能防止明显超配,不能替代数据库和应用层监控。


结语

异步数据库访问的核心不是语法,而是边界:

  • AsyncEngine 管理连接池和数据库连接资源;
  • AsyncSession 管理一个逻辑事务和 ORM 状态;
  • 一个 AsyncSession 不应被并发任务共享;
  • 事务把多条数据库操作绑定到提交边界;
  • 取消必须经过 finally 清理,并通常继续传播;
  • 超时不等于数据库端一定停止;
  • 并发一致性需要条件更新、锁、版本号和约束共同实现;
  • 重试必须建立在幂等和可重放之上;
  • 客户端请求生命周期不适合承载所有必须完成的业务操作。

当这些关系被明确之后,连接池配置、事务写法、任务调度和错误处理就不再是零散技巧,而是同一个系统模型的不同部分。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。