news 2026/8/25 5:20:50

第18章:FastAPI异步数据库访问与连接池

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
第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 框架中的灾难:

  1. 连接池耗尽pool_size太小 → 请求排队等连接;太大 → 数据库连接数过载,操作系统上下文切换开销剧增。
  2. N+1 查询:ORM 的惰性加载在循环中触发,100 行数据 = 101 次 SQL。数据库不是瓶颈——是你写的 ORM 用法是瓶颈。
  3. 阻塞事件循环:同步数据库驱动(psycopg2)在async def中调用——虽然 FastAPI 会把它扔到线程池,但线程池也有上限,500 并发就会有请求排队。
  4. 事务泄漏:某接口获取了数据库连接但忘记 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 时的表现
5P95 8000ms,大量 Connection Timeout
10P95 3500ms,偶尔 Timeout
20P95 600ms,稳定
50P95 500ms,稳定(但数据库 CPU 90%)
100P95 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),但可能导致笛卡尔积膨胀。

小白:“同步升级为异步,代码改动大吗?”

大师:"核心改动三点:

  1. 引擎:create_engine()create_async_engine()+asyncpg
  2. 会话:Session()AsyncSession()
  3. 查询: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 预加载,不会触发额外 SQL

4. 项目总结

优点 & 缺点对比

方案异步 SQLAlchemy + asyncpg同步 SQLAlchemy + psycopg2SQLModel (异步)raw asyncpg (无 ORM)
查询性能高(非阻塞 IO)中(线程池开销)最高
ORM 特性完整完整简化
连接池内置 + 可控参数内置 + 可控参数同左手动管理
N+1 解决selectinload/joinedload同左同左无此问题(手写 SQL)

适用场景

✓ 异步数据库适用于:

  1. 高并发读多写少的 API(电商商品列表、资讯 feed)
  2. 需要同时查询多个表的聚合接口(BFF 层)
  3. WebSocket 服务中的数据库访问
  4. 需要与多个外部 IO(HTTP API + DB + Redis)并发的场景
  5. PostgreSQL 数据库(asyncpg 支持最好)

✗ 不适用:

  1. SQLite 为主的项目——异步没有真正收益(aiosqlite 是假异步)
  2. 极简 CRUD(2-3 个表)——同步足够,异步增加心智负担

注意事项

  1. async with session.begin()vssession.commit()session.begin()自动管理事务的 begin/commit/rollback,强烈推荐。手动commit()+rollback()易遗漏。
  2. selectinloadvsjoinedloadselectinloadIN (...)查询,适合一对多和多对多。joinedload用 LEFT JOIN,适合一对一和多对一。错误选择会导致数据重复或性能倒退。
  3. pool_recycle设置:MySQL 默认 8 小时断开空闲连接。如果连接池中的连接超过 8 小时没使用,下次查询会报MySQL server has gone away。设pool_recycle=3600(1 小时)主动回收。
  4. 不要跨协程共享 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)——支持回滚到子事务,不影响外层事务。

思考题

  1. 初级:为订单查询接口增加 EXPLAIN ANALYZE 输出(PostgreSQL 的执行计划)。用db.execute(text("EXPLAIN ANALYZE SELECT ..."))查看是否使用了索引。

  2. 进阶:设计一个"读写分离"的数据库访问方案——写操作走主库,读操作走从库。如何在 FastAPI 的依赖注入中实现?提示:准备两个引擎engine_writeengine_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 实战修炼与源码剖析

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/25 5:20:09

从被动审核到主动风控:构建下一代视频内容安全体系

1. 从“救火队”到“防火墙”:视频内容安全范式的根本性转变如果你在2015年之前从事过视频平台的内容审核工作,那你大概率体验过什么叫“人海战术”和“事后诸葛亮”。那时候,我们面对的是海量的UGC(用户生成内容)&…

作者头像 李华
网站建设 2026/8/25 5:19:37

简历优化全攻略:提升求职成功率的实用技巧

1. 简历问题概述作为职场人士的第一张名片,简历质量直接决定了求职成功率。根据我多年担任面试官和职业咨询师的经验,90%的求职者都会在简历制作上犯各种错误。这些问题看似微小,却可能让你与心仪的工作失之交臂。2. 常见简历问题解析2.1 基础…

作者头像 李华
网站建设 2026/8/25 5:19:22

Java全栈开发工程师面试核心要点与实战策略

1. Java全栈开发工程师面试全景解析作为一位经历过数十场技术面试的Java全栈开发者,我深刻理解面试过程中的技术考察重点与应对策略。全栈开发岗位的特殊性在于,它要求候选人不仅要精通后端Java技术栈,还需要对前端开发、数据库设计、系统架构…

作者头像 李华
网站建设 2026/8/25 5:18:48

腾讯云直播音频审核实战:三种开启方式与避坑指南

1. 项目概述:为什么直播音频审核是刚需?做直播的朋友,尤其是游戏、秀场、电商带货这类UGC内容密集的领域,应该都遇到过类似的头疼事:主播一个不留神说了句不该说的,或者连麦的观众突然“口吐芬芳”&#xf…

作者头像 李华
网站建设 2026/8/25 5:16:57

2026固原工程建筑材料检测排名 TOP5 CMA 资质提供钢材检测、水泥检测、砂石检测 全覆盖联系方式推荐.txt

固原市区及周边建材检测机构星罗棋布,资质水平却参差不齐。建筑总包单位、建材生产厂家、市政工程项目与装修建设企业选材验收时,稍有不慎便会碰上无资质机构出具的检测报告,这类报告根本无法用于工程报审与竣工验收备案。小编实地走访、逐一…

作者头像 李华
网站建设 2026/8/25 5:16:27

人形机器人步频与储能技术:核心原理、优化方法与应用场景

这次我们来看一个近期在机器人领域引发广泛讨论的技术话题:人形机器人步频与储能。这并非一个具体的开源项目,而是一个聚焦于机器人核心运动性能与能源效率的关键技术方向。随着人形机器人从实验室走向更广阔的应用场景,其步行的“步频”和身…

作者头像 李华