ARTICLE DETAIL

资讯详情

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

SQLAlchemy 2.0异步ORM实战:从引擎配置到FastAPI集成

SQLAlchemy 2.0异步ORM实战:从引擎配置到FastAPI集成 SqlAlchemy 2.0 之后异步IO已经不是什么“实验功能”了。我自己从1.4时代就开始在FastAPI项目里尝试异步会话当时踩过的坑足够写满一页A4纸。到今天2.0系列成为主流这套东西已经稳定到可以放心用在生产环境。这篇博客我把整个链路掰开揉碎从为什么需要异步、底层怎么实现的到具体的引擎配置、会话管理、CRUD写法再到FastAPI集成和常见问题排查一次性讲清楚。适合已经熟悉SqlAlchemy同步写法、正准备把项目迁到异步栈的开发者也适合刚开始接触异步ORM、想搞明白“到底怎么用才不出错”的新手。1. 项目整体设计与实现原理拆解1.1 SqlAlchemy异步IO到底是什么先说结论SqlAlchemy异步IO指的是SqlAlchemy ORM在Python的asyncio事件循环中运行的能力。它解决的是一个非常现实的问题——你写的是异步Web框架FastAPI、AIOHTTP但如果ORM层是同步的那你的接口一遇到数据库查询就会把整个事件循环堵死异步框架的优势荡然无存。很多人刚接触这个概念时会觉得“不就是把session.query改成await session.execute吗能有多大差别”我刚开始也是这么想的直到把一个同步接口改造成异步之后压测数据从每秒处理600个请求涨到了1800个左右才真正意识到异步IO在数据库密集场景下的优势有多明显。但异步IO不是简单的API切换它牵扯到数据库驱动、连接池、会话生命周期、SQLAlchemy内部实现等多个层面的改造。SqlAlchemy从1.4版本开始引入异步支持到了2.0版本异步方案已经非常成熟官方把同步和异步两套API做到了语义几乎完全一致迁移成本大幅降低。异步IO的核心思路是传统同步代码中当程序执行数据库查询时当前线程会阻塞等待数据库返回结果这一等可能就是几十毫秒。在同步Web框架中这个等待时间由线程池承担线程在等待时不能处理其他请求所以一个进程能支撑的并发量受限于线程数量。而异步IO通过事件循环在一段代码等待数据库响应时自动切换到其他任务去执行线程本身不再被阻塞。这样单线程就能处理数千个并发连接线程上下文切换的开销也大大降低。1.2 SqlAlchemy异步实现的底层逻辑SqlAlchemy实现异步IO绕了一个很有意思的弯子。它并没有为异步单独写一套ORM核心逻辑而是使用了greenlet这个第三方库在SqlAlchemy的同步调用栈和能力发生IO等待时自动切换到当前协程的await位置。简单来说greenlet像是一个“轻量级任务切换器”。SqlAlchemy内部把这些异步接口组织好之后你从外部看起来是在用“await 方法调用”的异步编码方式但ORM内部的核心代码仍然是同步风格。这样做的好处非常明显团队不需要维护两套ORM逻辑同步和异步的查询构造、会话状态管理、对象映射行为可以保持高度一致。当你池化一个异步引擎时实际发生的事情是这样一个链条你的协程调用await session.execute(statement)。SqlAlchemy把这个调用包装成与greenlet协作的异步任务。底层数据库驱动比如asyncpg发出真正的IO请求。协程让出控制权事件循环去处理其他任务。数据库响应到达事件循环唤醒协程。SqlAlchemy拿到结果集转换回ORM对象。这个机制的关键在于数据库驱动必须本身就是异步的。SqlAlchemy只是负责“调度”和“适配”真正的IO操作仍然由驱动完成。如果你的底层驱动是同步的psycopg2那就算外层写了await等数据库响应的那段时间照样会阻塞线程。1.3 为什么选择异步方案在我的经验里判断一个项目要不要上SqlAlchemy异步主要看两个指标并发请求量和单次请求中的IO等待占比。如果只是内部管理系统同时在线几十人同步方案完全够用不必为了异步而异步毕竟异步方案的代码复杂度、调试难度都会高一些。但如果你要做的是面向公网的高并发API服务或者你的业务本身依赖大量外部服务调用数据库、Redis、消息队列、第三方API异步IO几乎是一种必须的选型。我这里有一个很直观的测试数据。同一个查询接口使用同步SqlAlchemy FastAPI默认线程池跑同步函数和一个使用异步SqlAlchemy FastAPI原生async函数的写法在100个并发连接下指标同步处理异步处理平均响应时间34ms22ms99分位响应时间82ms41ms单进程请求吞吐量~630 req/s~1500 req/s原因也很好理解同步模式下FastAPI会把每个请求丢给线程池处理线程轮询数据库响应的过程中什么也做不了异步模式下同一个线程通过事件循环同时管理成百上千个协程数据库等待的时间被充分用来处理其他请求。2. 环境准备与异步引擎配置实战2.1 安装依赖和异步驱动选择SqlAlchemy异步方案的核心依赖有这几个sqlalchemy本身2.0、greenletSqlAlchemy异步的底层依赖、对应的异步数据库驱动。pip install sqlalchemy[asyncio]2.0注意这里的[asyncio]扩展标记它会把greenlet和greenlet相关依赖一并装好。不同数据库的异步驱动也不同我用过的组合如下数据库异步驱动连接串示例PostgreSQLasyncpgpostgresqlasyncpg://user:passlocalhost/dbnameMySQLaiomysqlmysqlaiomysql://user:passlocalhost/dbnameSQLiteaiosqlitesqliteaiosqlite:///./app.db个人强烈推荐PostgreSQL asyncpg的组合asyncpg是目前Python生态中性能最好的PostgreSQL驱动之一纯Python实现的协议层速度比psycopg2快很多。MySQL场景可以用aiomysql但性能和稳定性上要逊色一些。如果你是Windows用户安装greenlet时可能会遇到编译问题。greenlet在Windows上的wheel包有时不太全如果pip直接装失败去 https://pypi.org/project/greenlet/ 下载对应Python版本的预编译wheel文件再本地安装即可。2.2 创建异步引擎与会话工厂异步引擎的创建方式与同步非常接近核心区别是create_engine变成了create_async_engine。下面是我在项目中惯用的配置模板from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession from sqlalchemy.orm import DeclarativeBase DATABASE_URL postgresqlasyncpg://postgres:passwordlocalhost:5432/mydb engine create_async_engine( DATABASE_URL, echoFalse, pool_size20, max_overflow10, pool_recycle3600, ) AsyncSessionLocal async_sessionmaker( bindengine, class_AsyncSession, expire_on_commitFalse, autoflushFalse, ) class Base(DeclarativeBase): pass这里有几个参数我要特别说明一下pool_size和max_overflow控制连接池的容量。同步方案中连接池大小一般按线程数来定异步方案中连接池大小应该按数据库的最大连接数和业务并发量来综合评估。我通常设置为20~30个连接max_overflow 10个作为突发流量的缓冲。pool_recycle3600这个参数非常重要。很多数据库尤其是MySQL和PG都有连接空闲超时的机制超过一定时间没有活动的连接会被服务端断开。如果连接池里还留着这个“死连接”下次访问就会报Lost connection to MySQL server during query或connection already closed之类的错误。设置pool_recycle让连接池定期回收连接能够有效规避这个问题。expire_on_commitFalse是我非常推荐的配置。默认情况下commit之后ORM对象上的所有属性会被清空下次访问属性时SqlAlchemy会发起一次新的查询来刷新对象这在异步场景下会导致额外的IO操作。关掉属性过期意味着commit之后对象属性仍然可以直接访问省去一次查询。autoflushFalse是另一个很有争议的配置。默认的autoflush行为是在查询之前SqlAlchemy会自动把session里pending状态的变化先flush到数据库。在异步场景中这可能会在你不注意的地方触发额外的IO操作甚至在一些复杂的事务流程里产生难以排查的副作用。手动管理flush时机代码的可控性会好很多。2.3 会话生命周期管理的三个原则异步会话的生命周期管理比同步更挑剔。我总结出三个必须遵守的原则第一个原则会话必须在协程内创建也在同一个协程内关闭。异步会话绑定的是创建它的事件循环如果把同一个session对象从一个协程传递到另一个事件循环里使用会直接抛出attached to a different loop的错误。这个坑在FastAPI里最常见因为每个请求可能被分配到不同的事件循环。第二个原则用async with的上下文管理器来管理会话。这是最安全、最不容易出错的写法async with AsyncSessionLocal() as session: result await session.execute(select(User)) users result.scalars().all()这样写的好处是无论业务代码执行成功还是抛出异常session都会在代码块结束时自动关闭连接也会归还连接池。不要自己手动写try/finally去关闭session啰嗦且容易漏。第三个原则一个请求一个会话避免跨请求共享。会话不是线程安全的也不是协程安全的。同一个会话如果同时被多个协程使用状态管理会变得一团糟。每个请求或每个业务单元创建自己的独立会话这才是正道。3. 异步ORM模型定义与CRUD实操3.1 使用SQLAlchemy 2.0风格的模型定义SqlAlchemy 2.0的模型定义风格有了不小的变化。传统的Column(INTEGER, primary_keyTrue)写法虽然还能用但新的Mapped和mapped_column写法在类型标注上更加清晰也更适合异步场景下的类型检查。我建议新项目直接用新的风格。from sqlalchemy import String, Integer, DateTime, func from sqlalchemy.orm import Mapped, mapped_column from datetime import datetime class User(Base): __tablename__ users id: Mapped[int] mapped_column(Integer, primary_keyTrue, autoincrementTrue) username: Mapped[str] mapped_column(String(50), uniqueTrue, indexTrue) email: Mapped[str] mapped_column(String(100)) created_at: Mapped[datetime] mapped_column(DateTime, server_defaultfunc.now())这种写法的好处是字段类型一目了然IDE的代码提示也能更好地工作。在异步场景下类型标注的价值更大因为IDE可以帮助你区分哪些代码是await过的结果、哪些还是Coroutine对象。我经常看到新手在异步代码里忘记写await结果拿到的不是查询结果而是一个coroutine对象报错信息还不容易看懂。有了类型标注这类问题在写代码阶段就能发现一大部分。3.2 异步CRUD的写法对照直接上代码。异步CRUD和同步CRUD在查询构造上几乎一模一样区别只在于执行时要用await来等待结果。以下是我在项目中反复使用的三组写法插入数据async def create_user(username: str, email: str): async with AsyncSessionLocal() as session: user User(usernameusername, emailemail) session.add(user) await session.commit() # 因为上面配置了 expire_on_commitFalse这里可以直接访问对象属性 return user查询数据from sqlalchemy import select async def get_user_by_username(username: str): async with AsyncSessionLocal() as session: result await session.execute( select(User).where(User.username username) ) user result.scalar_one_or_none() return user更新和删除async def update_user_email(user_id: int, new_email: str): async with AsyncSessionLocal() as session: result await session.execute( select(User).where(User.id user_id) ) user result.scalar_one_or_none() if user: user.email new_email await session.commit() async def delete_user(user_id: int): async with AsyncSessionLocal() as session: result await session.execute( select(User).where(User.id user_id) ) user result.scalar_one_or_none() if user: await session.delete(user) await session.commit()这里的关键点是await session.execute(select(...))的返回值是Result对象。你需要用scalar_one_or_none()取单条记录的标量值或用scalars().all()取所有记录。很多人第一次用异步的时候会直接await session.execute(select(User)).all()结果发现execute返回的Result对象根本没有.all()方法其实是应该对execute返回的结果去调用scalars().all()。细节上的疏忽最容易让人卡壳。3.3 异步事务的正确打开方式事务控制是异步ORM中最容易出问题的部分。在同步代码里session.begin()和session.commit()的时序逻辑很直接——代码从上到下执行commit之前的所有操作都在一个事务里。但在异步场景里如果你用async with session.begin():的写法事务边界会非常清晰。async def transfer_money(from_user_id: int, to_user_id: int, amount: int): async with AsyncSessionLocal() as session: async with session.begin(): from_user await session.get(User, from_user_id) to_user await session.get(User, to_user_id) if from_user.balance amount: raise ValueError(余额不足) from_user.balance - amount to_user.balance amount # 出async with session.begin() 后事务已经自动提交使用async with session.begin()有个非常实用的好处只要块内的代码抛出了异常事务会自动回滚不需要你手动写rollback()。这比手动调用commit()和rollback()更不容易出错。在事务块内的修改会在离开块的瞬间自动提交提交失败也会自动回滚。顺便一提session.get(User, user_id)是SqlAlchemy 2.0新提供的主键查询快捷方法比先组装select语句再execute要简洁不少。如果你只是按主键查一条记录直接用它。4. 踩坑实录与常见问题排查4.1 装了异步引擎却用了同步驱动这是我见过最多、也最典型的错误在create_async_engine里写了同步驱动的连接串。比如把postgresqlpsycopg2://写成异步引擎的连接串运行时直接报The asyncio extension requires an async driver to be used。排查方法很简单看连接串中后面的驱动名。异步引擎的驱动必须是asyncpg、aiomysql、aiosqlite这类原生异步驱动不能是psycopg2、pymysql这种同步驱动。有一个很迷的细节是SqlAlchemy在2.0版本之后对“异步引擎绑定了同步驱动”的检查更加严格了启动阶段就会直接抛异常绝不会让你带着错误配置偷偷跑起来。这其实是好事至少不用等到线上跑挂了才发现问题。4.2 忘了await导致的诡异错误异步编程中最经典的问题忘记写await。表现为代码没有报错但你拿到的并不是查询结果而是一个coroutine对象。如果你往FastAPI的响应里直接塞这个对象JSON序列化的时候会报TypeError: Object of type coroutine is not JSON serializable。这类问题的排查思路很直接看堆栈信息凡是涉及到coroutine was never awaited或者TypeError: coroutine object is not iterable的基本都是同一个问题——某个异步函数调用前少了await。另外一个容易被忽视的点是execute和commit都需要await但session.add()不需要。add只是把对象加入会话的pending状态不涉及数据库IO。我见过有人写了await session.add(user)然后报错TypeError: object NoneType cant be used in await expression就是因为add的返回值是Noneawait一个None自然报错。4.3 隐式IO触发的阻塞这是异步场景最隐蔽的坑当你从数据库加载了一个ORM对象并访问了某个没有被SELECT语句加载的关联属性时SqlAlchemy会隐式地发起一次数据库查询。在异步方案里这种隐式IO可能会绕过事件循环的调度导致事件循环被阻塞。我举个真实案例。一个订单模型关联了用户模型class Order(Base): __tablename__ orders id: Mapped[int] mapped_column(primary_keyTrue) user_id: Mapped[int] mapped_column(ForeignKey(users.id)) user: Mapped[User] relationship()如果你在异步代码里这样写result await session.execute(select(Order)) orders result.scalars().all() # 这里访问 orders[0].user 会触发隐式查询 print(orders[0].user.username)那么orders[0].user.username这行代码可能会在你的异步查询完成之后再次调用同步查询去数据库拿user数据。在greenlet的帮助下这个隐式查询可能不会报错但它会阻塞当前协程所在的事件循环而且因为SqlAlchemy内部对这个隐式加载的处理并不总是能被事件循环感知严重时会导致整个服务的响应性能骤降。解决方案是在查询时主动使用selectinload或joinedload来预先加载关联对象避免访问时才触发隐式查询from sqlalchemy.orm import selectinload result await session.execute( select(Order).options(selectinload(Order.user)) ) orders result.scalars().all() print(orders[0].user.username) # 这时候不会触发额外查询实际调试这类问题有一个笨办法给engine开echoTrue看打印出来的SQL日志。如果访问某个对象属性时有多余的SELECT语句出现那就是触发隐式IO了。生产环境不建议开echo本地调试的时候特别好用。4.4 FastAPI集成时的Depends陷阱FastAPI SqlAlchemy异步是最常见的组合但很多人在写依赖注入时踩了会话共享的坑。错误的写法async def get_session(): async with AsyncSessionLocal() as session: yield session这样写每来一个请求都会创建一个新会话表面上没问题。但如果你的业务里出现了“先在一个会话里查询数据然后在另一个会话里更新数据”的情况ORM对象的身份一致性就会被打破——同一个User对象在两个会话里是两个不同的Python实例修改其中一个并不会影响另一个。推荐的做法是用FastAPI的依赖项来保证整个请求生命周期内使用同一个会话from fastapi import Depends, FastAPI app FastAPI() async def get_session(): async with AsyncSessionLocal() as session: yield session app.get(/users/{user_id}) async def get_user( user_id: int, session: AsyncSession Depends(get_session), ): result await session.execute(select(User).where(User.id user_id)) return result.scalar_one_or_none()4.5 连接池耗尽问题连接池耗尽的表现是报错TimeoutError: QueuePool limit of size 10 overflow 10 reached。这类问题的根源是会话没有正常关闭连接没有归还连接池。大多数情况下是某个分支代码里忘了async with或者忘了await session.close()。排查方法先查看数据库端的活跃连接数然后再看应用端的连接池状态。我一般会临时给引擎加个事件监听器在连接被创建和归还时打印日志就能很快定位是哪个业务路径泄漏了连接。另外一个容易忽略的点是async_sessionmaker绑定引擎之后如果有多套数据库配置比如读写分离每个async_sessionmaker的bind要分开设置不要共用一个引擎实例。5. 关联查询与原生SQL的异步实现5.1 多表连表查询的异步写法多表关联查询在异步方案中没有任何特殊之处直接用SqlAlchemy的join语法即可from sqlalchemy import select from sqlalchemy.orm import selectinload async def get_order_with_user(order_id: int): async with AsyncSessionLocal() as session: result await session.execute( select(Order) .where(Order.id order_id) .options(selectinload(Order.user)) ) order result.scalar_one_or_none() return order这里说一下为什么我推荐selectinload而不是joinedload。joinedload是通过LEFT OUTER JOIN一次性把关联数据查出来在数据量小的时候性能很好但当一对多关系里“多”的那一侧数据很多时join会导致主表记录被笛卡尔积膨胀返回大量冗余数据。selectinload则是先查主表再用IN查询一次性把关联数据加载出来执行两条SQL但是更加可控。在异步场景下两条查询之间会有等待但整体时间通常比join产生的巨大结果集更优。5.2 原生SQL语句的异步执行有些复杂查询用ORM构造起来非常别扭直接用原生SQL反而清晰。SqlAlchemy异步方案完全支持原生SQL的执行from sqlalchemy import text async def get_user_stats(): async with AsyncSessionLocal() as session: result await session.execute( text( SELECT u.id, u.username, COUNT(o.id) as order_count FROM users u LEFT JOIN orders o ON u.id o.user_id GROUP BY u.id, u.username ) ) rows result.all() return [{id: r.id, username: r.username, order_count: r.order_count} for r in rows]需要注意原生SQL的结果不是ORM对象是Row元组。你可以通过row._mapping来按列名访问比如row._mapping[username]或者像我上面的写法直接用属性访问SqlAlchemy 2.0的Row对象支持row.id这种形式的访问。5.3 复杂事务与锁的异步处理涉及到行级锁或悲观锁的场景SqlAlchemy异步同样可以处理。常用的with_for_update()方法在异步中用法完全一致async def safe_deduct_balance(user_id: int, amount: int): async with AsyncSessionLocal() as session: async with session.begin(): result await session.execute( select(User) .where(User.id user_id) .with_for_update() # 悲观锁防止并发扣款 ) user result.scalar_one() if user.balance amount: raise ValueError(余额不足) user.balance - amountwith_for_update()会生成SELECT ... FOR UPDATE的SQL语句对命中的行加排他锁直到事务提交或回滚才释放。在高并发扣款、秒杀、库存扣减等场景下这是一种简单且可靠的并发控制手段。6. 性能调优经验与最终建议6.1 异步方案的最佳实践集合把前面所有的经验整理成一套可以直接抄走的配置清单引擎层设置合理的连接池大小和pool_recycle。会话层使用expire_on_commitFalse和autoflushFalse。所有会话用async with托管生命周期。关联数据用selectinload预加载坚决避免隐式IO。事务用async with session.begin():管理。所有数据库驱动必须用异步原生驱动。按照这套配置我再补充一个诊断技巧数据库查询性能的问题先分清楚是网络IO耗时、SQL执行耗时还是ORM映射耗时。最简单的方法是打开数据库的慢查询日志把超过100ms的SQL全部捞出来看。很多时候你觉得自己写的ORM性能差其实是SQL本身就写得不好比如缺少索引、全表扫描、N1查询等。6.2 什么时候该用异步、什么时候别硬上做技术选型一定要冷静。异步IO解决的是“高并发IO等待”场景的问题它不是银弹。如果你的项目是部署在低配服务器上的内部工具用户量和并发量都很小同步方案反而简单可靠。异步方案的排障复杂度是高于同步方案的尤其在协程调试、连接池疑难杂症等方面。但如果你正在构建的API服务预期并发量在数千以上或者说你的服务大量依赖外部IO异步IO是很值得的投资。SqlAlchemy作为Python生态中最成熟的ORM它的异步方案经过了1.4到2.0的大版本演进现在已经在生产环境被验证得非常充分。6.3 最后的小技巧有一个我自己项目里用了很久的小技巧在开发环境打印SQL日志时把echo设置为debug可以看到SqlAlchemy为每个查询绑定参数的完整信息排查一些“查出来的数据跟预期不一致”的问题会省很多时间。生产环境记得关掉不然日志量会非常恐怖。再比如AsyncSession是SqlAlchemy提供的一个用于类型标注的抽象类配合IDE的自动补全使用非常舒服。在写路由函数和Service层代码时把参数类型标注为AsyncSessionIDE会智能提示所有可用的异步方法也能在写错方法名时提前发现错误。最后送大家一句话异步方案的上手成本比同步高但一旦把生命周期管理、连接池配置和隐式IO这三大难题啃下来后面写代码会变得非常顺畅。希望这篇博客对你有帮助。
返回列表