Python 基础体系 · 第 85/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。
Python 异步数据库访问:连接池、事务、取消、并发和一致性
异步数据库访问的难点不在于把 def 改成 async def,而在于同时处理五个相互影响的问题:
- 连接池:有限的数据库连接如何被大量请求复用。
- 事务:多条 SQL 如何构成一个不可分割的业务操作。
- 取消:客户端断开、请求超时或服务关闭时,正在执行的数据库操作如何收尾。
- 并发:多个协程、多个请求和多个进程如何同时访问数据库。
- 一致性:在并发读写、重试和部分失败下,系统如何避免错误数据。
这五个问题不能分别孤立处理。例如,请求取消可能发生在事务提交前;连接池耗尽可能表现为接口超时;把同一个 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 文档将 Session、AsyncSession 和数据库连接都描述为有状态对象;一个会话代表一个逻辑事务,并且同一 AsyncSession 不适合被多个 asyncio 任务同时使用。(docs.sqlalchemy.org)
因此应区分三种“并发”:
| 层次 | 含义 | 是否能并行 |
|---|---|---|
| 协程并发 | 多个任务交错执行 | 可以 |
| 数据库连接并发 | 多个连接同时向数据库发请求 | 可以 |
| 同一连接上的命令 | 在同一事务连接上交错操作 | 通常不应并行 |
异步 API 解决的是“等待数据库时不阻塞事件循环”,不是“让一个数据库连接同时执行多条 SQL”。
二、连接池:控制数据库连接的并发入口
2.1 连接池解决什么问题
数据库连接不是普通的 Python 对象。建立连接通常需要:
- 创建 TCP 连接;
- 完成数据库协议握手;
- 认证;
- 设置会话参数;
- 可能执行初始化 SQL。
如果每个 HTTP 请求都新建连接,连接建立成本会重复发生,更重要的是,突发流量会直接把数据库连接数推高。
连接池把连接生命周期改成:
创建少量连接
↓
请求到来
↓
从池中借出连接
↓
执行事务和 SQL
↓
归还连接
连接池的容量不是“应用能处理的请求数”,而是“同一进程可以同时占用的数据库连接数”。
设:
W:应用进程或 worker 数量;P:每个进程的池大小;O:每个进程允许额外创建的溢出连接;C_db:数据库允许该应用使用的最大连接数。
应用理论上最多可能占用:
要避免应用自身突破数据库连接上限,至少需要满足:
其中 C_other 表示管理工具、迁移程序、监控组件和其他服务占用的连接。
例如:
- 4 个 worker;
- 每个 worker
pool_size=10; - 每个 worker
max_overflow=5; - 数据库为该服务预留 70 个连接;
- 其他组件需要 5 个连接。
则:
在这个简单估算下还有 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 转账的原子性
转账需要满足不变量:
两条更新必须同时发生,否则可能出现:
扣款成功
进程崩溃
入账没有执行
事务将两条更新绑定到同一个提交点:
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()
这段代码先提交了扣款,再调用外部服务,最终无法让数据库事务覆盖整个跨系统操作。更合理的设计通常是:
- 在一个数据库事务中创建支付意图或待处理记录;
- 提交;
- 调用外部服务;
- 通过幂等键和状态机记录结果;
- 失败时由补偿任务重试。
数据库事务不能跨越任意外部系统获得原子性。
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 文档特别提醒,TaskGroup 和 asyncio.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,说明其他事务已经修改了这条记录。应用可以:
- 重新读取;
- 重新计算;
- 重试有限次数;
- 或向调用者报告冲突。
乐观锁的核心公式是:
更新后:
任何并发事务只要仍然使用旧版本 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)
)
first 和 second 是否相同,不由 Python 的 async 关键字决定,而由:
- 数据库隔离级别;
- 查询是否使用同一事务;
- 其他事务何时提交;
- 数据库的 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,除非确实需要特殊资源管理。
当出现连接池耗尽时,应检查:
- 是否存在未关闭的
AsyncSession; - 是否有事务打开后长时间等待外部 I/O;
- 是否把数据库查询结果交给了长生命周期任务;
- 是否在应用启动时重复创建 Engine;
- worker 数量乘以池大小是否超过数据库限制;
- 是否有慢查询或锁等待,使连接长时间无法归还;
- 是否发生了取消异常但清理代码没有执行。
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
关键点有三个:
- 取消发生在任意
await点,可能发生在获取连接、执行 SQL、读取结果或提交阶段; - 回滚不是业务成功的反义词,而是事务清理动作;
- 连接归还必须发生,即使请求没有正常返回响应。
取消处理代码应保持短小:
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 ...
诊断顺序:
- 统计请求进入数据库层的数量;
- 查看连接池已签出连接数;
- 检查事务持续时间;
- 查看数据库锁等待和慢查询;
- 检查取消路径是否释放会话;
- 按 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
第五步:确定取消语义
问清楚:
客户端断开后,这个操作还能不能停止?
如果不能停止,就不要把它绑定在请求任务上。
第六步:确定重试语义
问清楚:
如果提交结果未知,再执行一次会发生什么?
如果可能重复扣款、重复发货或重复创建,就必须先加入幂等键和状态查询。
第七步:计算连接上限
使用:
然后再结合平均事务时间、慢查询、锁等待和峰值并发进行压测。公式只能防止明显超配,不能替代数据库和应用层监控。
结语
异步数据库访问的核心不是语法,而是边界:
AsyncEngine管理连接池和数据库连接资源;AsyncSession管理一个逻辑事务和 ORM 状态;- 一个
AsyncSession不应被并发任务共享; - 事务把多条数据库操作绑定到提交边界;
- 取消必须经过
finally清理,并通常继续传播; - 超时不等于数据库端一定停止;
- 并发一致性需要条件更新、锁、版本号和约束共同实现;
- 重试必须建立在幂等和可重放之上;
- 客户端请求生命周期不适合承载所有必须完成的业务操作。
当这些关系被明确之后,连接池配置、事务写法、任务调度和错误处理就不再是零散技巧,而是同一个系统模型的不同部分。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Alembic 数据库迁移:Revision、自动生成、分支、上线和回滚
- 下一篇:Python 与 Redis:连接池、Pipeline、事务、缓存和分布式锁
- 延伸:SQLAlchemy 2.0:Engine、Session、映射、查询、事务和 N+1
- 延伸:Python asyncio 完整基础:事件循环、协程、Future 与调度
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论