ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

第18章:FastAPI异步数据库访问与连接池 1. 项目背景业务场景第 16 章的任务协作 API 中数据库访问用的是同步 SQLAlchemy def端点。当并发量上去后性能瓶颈暴露了订单查询接口在 500 并发下 P95 延迟达到4200ms其中 80% 的时间在等待数据库连接。数据库连接池配的是默认pool_size5, max_overflow10——高峰时 500 个请求争抢 15 个连接排队时间比查询时间还长。有一个热门商品详情接口——一个请求要查商品表、评论表、用户表总共 3 次数据库。更糟的是获取每条评论的作者名时由于 ORM 的惰性加载lazy loading又额外触发了一次查询。100 条评论 101 次数据库查询——这就是著名的N1 问题。运维说“你们服务把 PostgreSQL 的连接数跑满了——500 个连接。但数据库服务器只有 4 核最佳连接数应该是2 × CPU 1 ≈ 9。”痛点同步数据库访问在异步 Web 框架中的灾难连接池耗尽pool_size太小 → 请求排队等连接太大 → 数据库连接数过载操作系统上下文切换开销剧增。N1 查询ORM 的惰性加载在循环中触发100 行数据 101 次 SQL。数据库不是瓶颈——是你写的 ORM 用法是瓶颈。阻塞事件循环同步数据库驱动psycopg2在async def中调用——虽然 FastAPI 会把它扔到线程池但线程池也有上限500 并发就会有请求排队。事务泄漏某接口获取了数据库连接但忘记 commit/rollback连接一直处于idle in transaction状态锁住行数据其他接口的 UPDATE 全部挂起。本章将同步数据库访问升级为异步用AsyncSessionasyncpg结合连接池调优解决上述所有瓶颈。2. 项目设计场景监控屏上数据库连接数曲线像过山车——从 5 飙到 200 再掉回 3。大师把 DBA 也叫来了。小胖指着监控屏“这台数据库服务器的连接数怎么跟股票似的一会儿 5 个一会儿 200 个”DBA 老王“你们写代码的时候管过连接池吗ORM 默认pool_size5高峰期不够用就每个人自己开连接——炸了。数据库连接是重型资源——一个连接在 PostgreSQL 里就是一个操作系统进程fork 模型200 个连接就是 200 个进程抢 4 个核。”大师“老王的痛点是第一手经验。我们先建立连接池的直觉”技术映射数据库连接池 预创建一批连接并复用。就像银行柜台——开 5 个窗口pool_size正常情况下够用。业务高峰期允许再临时开 10 个max_overflow总共 15 个窗口。高峰期过后临时窗口关闭回收。pool_timeout是如果所有窗口都忙客人最多等多久默认 30 秒。小白“那到底配多少个连接才合适550500”大师“有一个经验公式——(2 × CPU 核心数) 有效磁盘数。但更重要的是实测——通过压测找到最优值。因为连接池大小和 QPS、查询复杂度、网络延迟都相关。没有万能公式只有压测数据。”pool_size并发 500 时的表现5P95 8000ms大量 Connection Timeout10P95 3500ms偶尔 Timeout20P95 600ms稳定50P95 500ms稳定但数据库 CPU 90%100P95 700ms变差上下文切换开销 连接收益“看到没——50 比 20 提升不大100 反而倒退。这就是连接池的’最优区间’。”小胖“那 N1 问题呢我代码里确实写了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()。”# ❌ N11 N 条 SQLtaskssession.execute(select(Task)).scalars().all()fortaskintasks:forcommentintask.comments:# 每条评论触发一次 SQLprint(comment.author.username)# 又触发 SQL# ✓ 预加载1 条 SQLJOIN 所有关联stmt(select(Task).options(selectinload(Task.comments).selectinload(TaskComment.author)))taskssession.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. 项目实战——异步数据库升级与连接池调优环境准备pipinstallsqlalchemy2.0.36asyncpg0.30.0aiosqlite0.20.0# 生产: asyncpg (PostgreSQL); 开发: aiosqlite (SQLite)分步实现步骤一创建异步数据库引擎和会话工厂目标异步驱动 连接池配置app/core/database.pyfromsqlalchemy.ext.asyncioimport(create_async_engine,AsyncSession,async_sessionmaker,AsyncEngine,)fromapp.core.configimportsettingsdefcreate_engine()-AsyncEngine:创建异步数据库引擎returncreate_async_engine(settings.DATABASE_URL,# postgresqlasyncpg://user:passhost:5432/dbechosettings.DEBUG,# ═══════ 连接池调优参数 ═══════pool_size20,# 常驻连接数max_overflow10,# 额外允许超出 pool_size 的连接数高峰弹性pool_timeout30,# 等待可用连接的超时秒数超时抛 QueuePool 错误pool_recycle3600,# 连接最大存活秒数防 MySQL 8 小时超时pool_pre_pingTrue,# 使用前先检测连接是否存活防断连# ═══════ 可选自定义连接参数 ═══════connect_args{timeout:10,# asyncpg 连接超时command_timeout:30,# asyncpg 单条 SQL 超时},)# 异步会话工厂AsyncSessionLocalasync_sessionmaker(bindNone,# 运行时动态绑定class_AsyncSession,expire_on_commitFalse,# 提交后不过期对象避免 DetachedInstanceErrorautoflushFalse,)asyncdefget_db()-AsyncSession:# type: ignoreFastAPI 异步数据库依赖——每个请求独立的 AsyncSessionasyncwithAsyncSessionLocal(bindcreate_engine())assession:try:yieldsessionawaitsession.commit()exceptException:awaitsession.rollback()raise# async with 自动调用 session.close()步骤二改造 Repository 为异步目标最小改动所有查询加 awaitapp/domains/order/repository.pyfromsqlalchemyimportselect,funcfromsqlalchemy.ext.asyncioimportAsyncSessionfromsqlalchemy.ormimportselectinloadfromapp.models.orderimportOrderfromapp.models.userimportUserclassOrderRepository:异步订单仓库asyncdeffind_by_id_with_user(self,db:AsyncSession,order_id:int)-Order|None:查询订单 预加载用户信息避免 N1stmt(select(Order).options(selectinload(Order.user))# 预加载关联 User.where(Order.idorder_id))resultawaitdb.execute(stmt)returnresult.scalar_one_or_none()asyncdefsearch(self,db:AsyncSession,user_id:int|NoneNone,status:str|NoneNone,page:int1,size:int20,)-tuple[list[Order],int]:搜索订单——异步分页查询stmtselect(Order)ifuser_idisnotNone:stmtstmt.where(Order.user_iduser_id)ifstatusisnotNone:stmtstmt.where(Order.statusstatus)# 总数count_stmtselect(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))itemslist((awaitdb.execute(stmt)).unique().scalars().all())returnitems,totalasyncdefcreate(self,db:AsyncSession,data:dict)-Order:orderOrder(**data)db.add(order)awaitdb.flush()# 异步 flushreturnorder关键点所有db.execute()前面加await所有方法声明为async def。业务逻辑不需要改动——只是加了 async/await 标注。步骤三改造 Service 层目标异步编排保留事务边界app/domains/order/service.pyclassOrderService:异步订单服务def__init__(self,repo:OrderRepository|NoneNone):self.reporepoorOrderRepository()asyncdefcreate_order(self,db:AsyncSession,user_id:int,data:dict)-Order:创建订单——异步事务# 事务边界使用 db.begin()asyncwithdb.begin():# 在同一个事务中执行data[user_id]user_id data[status]pendingorderawaitself.repo.create(db,data)returnorderasyncdeflist_orders(self,db:AsyncSession,user_id:int,status:str|None,page:int,size:int)-dict:异步查询订单列表items,totalawaitself.repo.search(db,user_id,status,page,size)return{items:[self._to_dict(o)foroinitems],total:total,page:page,size:size,}步骤四创建测试对比脚本目标量化异步升级的性能提升scripts/benchmark_orders.pyimporttimeimportasynciofromapp.domains.order.serviceimportOrderServicefromapp.core.databaseimportget_dbasyncdefbench_async():异步版本性能测试db_genget_db()dbawaitanext(db_gen)# 获取 AsyncSessionserviceOrderService()starttime.perf_counter()tasks[service.list_orders(db,user_idi%100,statusNone,page1,size20)foriinrange(500)]resultsawaitasyncio.gather(*tasks)elapsedtime.perf_counter()-startprint(fAsync 500 requests:{elapsed:.2f}s ({500/elapsed:.0f}req/s))asyncio.run(bench_async())步骤五启动服务并验证# 启动 PostgreSQLDockerdockerrun-d--namepg-test-ePOSTGRES_PASSWORDtest123\-p5432:5432 postgres:16-alpine# 设置环境变量$env:DATABASE_URLpostgresqlasyncpg://postgres:test123localhost:5432/testdb# 执行迁移alembic upgradehead# 启动服务uvicorn app.main:app--reload# 测试异步查询接口curl-shttp://localhost:8000/api/v1/orders?page1size20|python-mjson.tool可能遇到的坑asyncpgvspsycopg3都是异步 PostgreSQL 驱动。asyncpg性能极高纯 Python 异步协议实现但不支持某些复杂类型如自定义 composite type。psycopg3 功能更全面但性能略低。生产环境推荐asyncpg。SQLite 不支持异步aiosqlite在底层其实是线程池模拟异步——它仍然会阻塞线程。开发环境可以用但不要测异步性能。expire_on_commitFalse异步环境下这个设置尤其重要——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.repositoryimportOrderRepositorypytest.mark.asyncioasyncdeftest_async_create_and_query():enginecreate_engine()asyncwithAsyncSessionLocal(bindengine)asdb:repoOrderRepository()orderawaitrepo.create(db,{user_id:1,product_name:Test,quantity:1,unit_price:10.0,total_amount:10.0,})assertorder.idisnotNone# 查询无 N1foundawaitrepo.find_by_id_with_user(db,order.id)assertfound.userisnotNone# selectinload 预加载不会触发额外 SQL4. 项目总结优点 缺点对比方案异步 SQLAlchemy asyncpg同步 SQLAlchemy psycopg2SQLModel (异步)raw asyncpg (无 ORM)查询性能高非阻塞 IO中线程池开销高最高ORM 特性完整完整简化无连接池内置 可控参数内置 可控参数同左手动管理N1 解决selectinload/joinedload同左同左无此问题手写 SQL适用场景✓ 异步数据库适用于高并发读多写少的 API电商商品列表、资讯 feed需要同时查询多个表的聚合接口BFF 层WebSocket 服务中的数据库访问需要与多个外部 IOHTTP API DB Redis并发的场景PostgreSQL 数据库asyncpg 支持最好✗ 不适用SQLite 为主的项目——异步没有真正收益aiosqlite 是假异步极简 CRUD2-3 个表——同步足够异步增加心智负担注意事项async with session.begin()vssession.commit()session.begin()自动管理事务的 begin/commit/rollback强烈推荐。手动commit()rollback()易遗漏。selectinloadvsjoinedloadselectinload用IN (...)查询适合一对多和多对多。joinedload用 LEFT JOIN适合一对一和多对一。错误选择会导致数据重复或性能倒退。pool_recycle设置MySQL 默认 8 小时断开空闲连接。如果连接池中的连接超过 8 小时没使用下次查询会报MySQL server has gone away。设pool_recycle36001 小时主动回收。不要跨协程共享 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_pingTruepool_recycle做防御添加 SQLAlchemy 的echo_poolTrue日志追踪连接生命周期。案例三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 层添加debugTrue参数开发环境输出 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 实战修炼与源码剖析
返回列表