第18章:FastAPI异步数据库访问与连接池
1. 项目背景
业务场景
第 16 章的任务协作 API 中,数据库访问用的是同步 SQLAlchemy +def端点。当并发量上去后,性能瓶颈暴露了:
- 订单查询接口在 500 并发下 P95 延迟达到4200ms,其中 80% 的时间在等待数据库连接。
- 数据库连接池配的是默认
pool_size=5, max_overflow=10——高峰时 500 个请求争抢 15 个连接,排队时间比查询时间还长。 - 有一个"热门商品详情"接口——一个请求要查商品表、评论表、用户表总共 3 次数据库。更糟的是,获取每条评论的作者名时,由于 ORM 的惰性加载(lazy loading),又额外触发了一次查询。100 条评论 = 101 次数据库查询——这就是著名的N+1 问题。
- 运维说:“你们服务把 PostgreSQL 的连接数跑满了——500 个连接。但数据库服务器只有 4 核,最佳连接数应该是
2 × CPU + 1 ≈ 9。”
痛点
同步数据库访问在异步 Web 框架中的灾难:
- 连接池耗尽:
pool_size太小 → 请求排队等连接;太大 → 数据库连接数过载,操作系统上下文切换开销剧增。 - N+1 查询:ORM 的惰性加载在循环中触发,100 行数据 = 101 次 SQL。数据库不是瓶颈——是你写的 ORM 用法是瓶颈。
- 阻塞事件循环:同步数据库驱动(psycopg2)在
async def中调用——虽然 FastAPI 会把它扔到线程池,但线程池也有上限,500 并发就会有请求排队。 - 事务泄漏:某接口获取了数据库连接但忘记 commit/rollback,连接一直处于"idle in transaction"状态,锁住行数据,其他接口的 UPDATE 全部挂起。
本章将同步数据库访问升级为异步,用AsyncSession+asyncpg结合连接池调优,解决上述所有瓶颈。
2. 项目设计
场景:监控屏上,数据库连接数曲线像过山车——从 5 飙到 200 再掉回 3。大师把 DBA 也叫来了。
小胖:(指着监控屏)“这台数据库服务器的连接数怎么跟股票似的?一会儿 5 个,一会儿 200 个?”
DBA 老王:“你们写代码的时候管过连接池吗?ORM 默认pool_size=5,高峰期不够用就每个人自己开连接——炸了。数据库连接是重型资源——一个连接在 PostgreSQL 里就是一个操作系统进程(fork 模型),200 个连接就是 200 个进程抢 4 个核。”
大师:“老王的痛点是第一手经验。我们先建立连接池的直觉:”
技术映射:数据库连接池 = 预创建一批连接并复用。就像银行柜台——开 5 个窗口(
pool_size),正常情况下够用。业务高峰期允许再临时开 10 个(max_overflow),总共 15 个窗口。高峰期过后,临时窗口关闭回收。pool_timeout是如果所有窗口都忙,客人最多等多久(默认 30 秒)。
小白:“那到底配多少个连接才合适?5?50?500?”
大师:“有一个经验公式——(2 × CPU 核心数) + 有效磁盘数。但更重要的是实测——通过压测找到最优值。因为连接池大小和 QPS、查询复杂度、网络延迟都相关。没有万能公式,只有压测数据。”
| pool_size | 并发 500 时的表现 |
|---|---|
| 5 | P95 8000ms,大量 Connection Timeout |
| 10 | P95 3500ms,偶尔 Timeout |
| 20 | P95 600ms,稳定 |
| 50 | P95 500ms,稳定(但数据库 CPU 90%) |
| 100 | P95 700ms,变差(上下文切换开销 > 连接收益) |
“看到没——50 比 20 提升不大,100 反而倒退。这就是连接池的’最优区间’。”
小胖:“那 N+1 问题呢?我代码里确实写了for comment in task.comments: print(comment.author.username)——这有什么问题?”
大师:“task.comments是 ORM 的 relationship。如果你没有预加载(eager loading),当你访问.comments时发一条 SQL,循环里每次访问.author又各发一条 SQL。100 个评论 = 1(查评论)+ 100(查每个作者)= 101 条 SQL。解决就一行:selectinload()。”
# ❌ N+1:1 + N 条 SQLtasks=session.execute(select(Task)).scalars().all()fortaskintasks:forcommentintask.comments:# 每条评论触发一次 SQL!print(comment.author.username)# 又触发 SQL!# ✓ 预加载:1 条 SQL(JOIN 所有关联)stmt=(select(Task).options(selectinload(Task.comments).selectinload(TaskComment.author)))tasks=session.execute(stmt).unique().scalars().all()技术映射:
selectinload是一种预加载策略——在一条 SQL 中用IN (task_ids)批量加载关联对象,存入 ORM 的 identity map。后续访问.comments直接从内存取,不发 SQL。joinedload是另一种(用 JOIN),但可能导致笛卡尔积膨胀。
小白:“同步升级为异步,代码改动大吗?”
大师:"核心改动三点:
- 引擎:
create_engine()→create_async_engine()+asyncpg - 会话:
Session()→AsyncSession() - 查询:
session.execute()→await session.execute()
业务层代码变化很小——把db.execute()前面加await就行。但依赖注入的管理方式要变:从yield变成async with。"
3. 项目实战——异步数据库升级与连接池调优
环境准备
pipinstallsqlalchemy==2.0.36asyncpg==0.30.0aiosqlite==0.20.0# 生产: asyncpg (PostgreSQL); 开发: aiosqlite (SQLite)分步实现
步骤一:创建异步数据库引擎和会话工厂(目标:异步驱动 + 连接池配置)
app/core/database.py:
fromsqlalchemy.ext.asyncioimport(create_async_engine,AsyncSession,async_sessionmaker,AsyncEngine,)fromapp.core.configimportsettingsdefcreate_engine()->AsyncEngine:"""创建异步数据库引擎"""returncreate_async_engine(settings.DATABASE_URL,# postgresql+asyncpg://user:pass@host:5432/dbecho=settings.DEBUG,# ═══════ 连接池调优参数 ═══════pool_size=20,# 常驻连接数max_overflow=10,# 额外允许超出 pool_size 的连接数(高峰弹性)pool_timeout=30,# 等待可用连接的超时秒数(超时抛 QueuePool 错误)pool_recycle=3600,# 连接最大存活秒数(防 MySQL 8 小时超时)pool_pre_ping=True,# 使用前先检测连接是否存活(防断连)# ═══════ 可选:自定义连接参数 ═══════connect_args={"timeout":10,# asyncpg 连接超时"command_timeout":30,# asyncpg 单条 SQL 超时},)# 异步会话工厂AsyncSessionLocal=async_sessionmaker(bind=None,# 运行时动态绑定class_=AsyncSession,expire_on_commit=False,# 提交后不过期对象(避免 DetachedInstanceError)autoflush=False,)asyncdefget_db()->AsyncSession:# type: ignore"""FastAPI 异步数据库依赖——每个请求独立的 AsyncSession"""asyncwithAsyncSessionLocal(bind=create_engine())assession:try:yieldsessionawaitsession.commit()exceptException:awaitsession.rollback()raise# async with 自动调用 session.close()步骤二:改造 Repository 为异步(目标:最小改动,所有查询加 await)
app/domains/order/repository.py:
fromsqlalchemyimportselect,funcfromsqlalchemy.ext.asyncioimportAsyncSessionfromsqlalchemy.ormimportselectinloadfromapp.models.orderimportOrderfromapp.models.userimportUserclassOrderRepository:"""异步订单仓库"""asyncdeffind_by_id_with_user(self,db:AsyncSession,order_id:int)->Order|None:"""查询订单 + 预加载用户信息(避免 N+1)"""stmt=(select(Order).options(selectinload(Order.user))# 预加载关联 User.where(Order.id==order_id))result=awaitdb.execute(stmt)returnresult.scalar_one_or_none()asyncdefsearch(self,db:AsyncSession,user_id:int|None=None,status:str|None=None,page:int=1,size:int=20,)->tuple[list[Order],int]:"""搜索订单——异步分页查询"""stmt=select(Order)ifuser_idisnotNone:stmt=stmt.where(Order.user_id==user_id)ifstatusisnotNone:stmt=stmt.where(Order.status==status)# 总数count_stmt=select(func.count()).select_from(stmt.subquery())total=(awaitdb.execute(count_stmt)).scalar()or0# 分页 + 预加载stmt=(stmt.options(selectinload(Order.user)).order_by(Order.created_at.desc()).offset((page-1)*size).limit(size))items=list((awaitdb.execute(stmt)).unique().scalars().all())returnitems,totalasyncdefcreate(self,db:AsyncSession,data:dict)->Order:order=Order(**data)db.add(order)awaitdb.flush()# 异步 flushreturnorder关键点:所有db.execute()前面加await,所有方法声明为async def。业务逻辑不需要改动——只是加了 async/await 标注。
步骤三:改造 Service 层(目标:异步编排保留事务边界)
app/domains/order/service.py:
classOrderService:"""异步订单服务"""def__init__(self,repo:OrderRepository|None=None):self.repo=repoorOrderRepository()asyncdefcreate_order(self,db:AsyncSession,user_id:int,data:dict)->Order:"""创建订单——异步事务"""# 事务边界:使用 db.begin()asyncwithdb.begin():# 在同一个事务中执行data["user_id"]=user_id data["status"]="pending"order=awaitself.repo.create(db,data)returnorderasyncdeflist_orders(self,db:AsyncSession,user_id:int,status:str|None,page:int,size:int)->dict:"""异步查询订单列表"""items,total=awaitself.repo.search(db,user_id,status,page,size)return{"items":[self._to_dict(o)foroinitems],"total":total,"page":page,"size":size,}步骤四:创建测试对比脚本(目标:量化异步升级的性能提升)
scripts/benchmark_orders.py:
importtimeimportasynciofromapp.domains.order.serviceimportOrderServicefromapp.core.databaseimportget_dbasyncdefbench_async():"""异步版本性能测试"""db_gen=get_db()db=awaitanext(db_gen)# 获取 AsyncSessionservice=OrderService()start=time.perf_counter()tasks=[service.list_orders(db,user_id=i%100,status=None,page=1,size=20)foriinrange(500)]results=awaitasyncio.gather(*tasks)elapsed=time.perf_counter()-startprint(f"Async 500 requests:{elapsed:.2f}s ({500/elapsed:.0f}req/s)")asyncio.run(bench_async())步骤五:启动服务并验证
# 启动 PostgreSQL(Docker)dockerrun-d--namepg-test-ePOSTGRES_PASSWORD=test123\-p5432:5432 postgres:16-alpine# 设置环境变量$env:DATABASE_URL="postgresql+asyncpg://postgres:test123@localhost:5432/testdb"# 执行迁移alembic upgradehead# 启动服务uvicorn app.main:app--reload# 测试异步查询接口curl-s"http://localhost:8000/api/v1/orders?page=1&size=20"|python-mjson.tool可能遇到的坑:
asyncpgvspsycopg3:都是异步 PostgreSQL 驱动。asyncpg性能极高(纯 Python 异步协议实现),但不支持某些复杂类型(如自定义 composite type)。psycopg3 功能更全面但性能略低。生产环境推荐asyncpg。- SQLite 不支持异步:
aiosqlite在底层其实是线程池模拟异步——它仍然会阻塞线程。开发环境可以用,但不要测异步性能。 expire_on_commit=False:异步环境下这个设置尤其重要——commit 后如果不 expire,后续访问 ORM 对象的属性不会触发惰性加载(因为 Session 可能已经关闭了)。
完整代码清单
本章完整代码见column/code/chapter18/,主要文件:
app/core/database.py:异步引擎 + 连接池配置app/domains/order/repository.py:异步 Repository(含selectinload预加载)app/domains/order/service.py:异步 Service(含async with db.begin()事务)scripts/benchmark_orders.py:性能对比脚本
测试验证
# tests/test_async_repo.pyimportpytestfromapp.core.databaseimportcreate_engine,AsyncSessionLocalfromapp.domains.order.repositoryimportOrderRepository@pytest.mark.asyncioasyncdeftest_async_create_and_query():engine=create_engine()asyncwithAsyncSessionLocal(bind=engine)asdb:repo=OrderRepository()order=awaitrepo.create(db,{"user_id":1,"product_name":"Test","quantity":1,"unit_price":10.0,"total_amount":10.0,})assertorder.idisnotNone# 查询(无 N+1)found=awaitrepo.find_by_id_with_user(db,order.id)assertfound.userisnotNone# selectinload 预加载,不会触发额外 SQL4. 项目总结
优点 & 缺点对比
| 方案 | 异步 SQLAlchemy + asyncpg | 同步 SQLAlchemy + psycopg2 | SQLModel (异步) | raw asyncpg (无 ORM) |
|---|---|---|---|---|
| 查询性能 | 高(非阻塞 IO) | 中(线程池开销) | 高 | 最高 |
| ORM 特性 | 完整 | 完整 | 简化 | 无 |
| 连接池 | 内置 + 可控参数 | 内置 + 可控参数 | 同左 | 手动管理 |
| N+1 解决 | selectinload/joinedload | 同左 | 同左 | 无此问题(手写 SQL) |
适用场景
✓ 异步数据库适用于:
- 高并发读多写少的 API(电商商品列表、资讯 feed)
- 需要同时查询多个表的聚合接口(BFF 层)
- WebSocket 服务中的数据库访问
- 需要与多个外部 IO(HTTP API + DB + Redis)并发的场景
- PostgreSQL 数据库(asyncpg 支持最好)
✗ 不适用:
- SQLite 为主的项目——异步没有真正收益(aiosqlite 是假异步)
- 极简 CRUD(2-3 个表)——同步足够,异步增加心智负担
注意事项
async with session.begin()vssession.commit():session.begin()自动管理事务的 begin/commit/rollback,强烈推荐。手动commit()+rollback()易遗漏。selectinloadvsjoinedload:selectinload用IN (...)查询,适合一对多和多对多。joinedload用 LEFT JOIN,适合一对一和多对一。错误选择会导致数据重复或性能倒退。pool_recycle设置:MySQL 默认 8 小时断开空闲连接。如果连接池中的连接超过 8 小时没使用,下次查询会报MySQL server has gone away。设pool_recycle=3600(1 小时)主动回收。- 不要跨协程共享 AsyncSession:同一个 AsyncSession 不应在多个协程中交替使用——它是单线程模型异步的,但内部状态不是协程安全的。
常见踩坑经验
案例一:MissingGreenlet错误
- 现象:
sqlalchemy.exc.MissingGreenlet: greenlet_spawn has not been called。 - 根因:在异步环境中访问了 ORM 对象的惰性加载属性,但当前不在数据库 Session 的上下文中。
- 解决:使用
selectinload()预加载需要的关联,或在 Session 关闭前访问完所有需要的属性。
案例二:连接池泄漏(QueuePool limit reached)
- 现象:服务运行几小时后所有请求返回
TimeoutError: QueuePool limit of size 20 overflow 10 reached。 - 根因:某个接口获取了 Session 但没有关闭——可能是缺少
await session.close(),或忘记在async with块中使用。 - 解决:排查所有
get_db()调用的地方;用pool_pre_ping=True+pool_recycle做防御;添加 SQLAlchemy 的echo_pool=True日志追踪连接生命周期。
案例三:async with db.begin()嵌套事务
- 现象:在已开启事务的
db中再次async with db.begin()——抛出InvalidRequestError: A transaction is already begun。 - 根因:SQLAlchemy 的 Session 不支持嵌套事务(非保存点)。
begin()只能调用一次。 - 解决:使用
db.begin_nested()开启保存点(savepoint)——支持回滚到子事务,不影响外层事务。
思考题
初级:为订单查询接口增加 EXPLAIN ANALYZE 输出(PostgreSQL 的执行计划)。用
db.execute(text("EXPLAIN ANALYZE SELECT ..."))查看是否使用了索引。进阶:设计一个"读写分离"的数据库访问方案——写操作走主库,读操作走从库。如何在 FastAPI 的依赖注入中实现?提示:准备两个引擎
engine_write和engine_read,在路由依赖中根据 HTTP 方法路由到不同的引擎。
答案提示:第 1 题在 Repository 层添加
debug=True参数,开发环境输出 EXPLAIN 结果。第 2 题的核心是自定义get_db(method: str)依赖——POST/PUT/DELETE → get_write_db(),GET → get_read_db()。但要注意"主从延迟"——刚写入的数据从库可能还没同步。第 25 章和第 28 章继续深入。
延伸阅读与资源
NumPy 从入门到生产落地:全链路实战指南(科学计算/向量化)
Redis 8 实战精讲:从 CRUD 到源码,构建高可用缓存系统
Redis 实战修炼与原理进阶
Python 3实战精进:从脚本到高并发订单引擎
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
MongoDB 实战进阶与内核修炼
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析
