ARTICLE DETAIL

资讯详情

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

FastAPI数据层实战:异步SQLAlchemy、CRUD与事务控制全解析

FastAPI数据层实战:异步SQLAlchemy、CRUD与事务控制全解析 很多朋友问到FastAPI的数据操作到底怎么写才是规范、能扛住生产环境的这期教程正好填上这个坑。前六篇我们把FastAPI的接口定义、参数校验、依赖注入和中间件都过了一遍但很多项目走到数据库这一层就卡住了——SQLAlchemy的同步写法放进async接口里容易阻塞事件循环ORM模型和Pydantic模型来回转换容易绕晕事务一多就不知道什么时候commit、什么时候rollback。这篇教程就把数据操作这一整条链路讲透从异步引擎配置、模型设计、CRUD接口到批量导入导出和事务控制全部给出可以直接抄进项目里的写法。内容偏向实战会结合我实际项目里踩过的坑来说明适合已经会写基础FastAPI接口、想正经做数据层的开发者参考也适合正在准备FastAPI相关面试的人用来补全知识点。1. 数据操作的整体设计与方案选型1.1 为什么数据层必须单独设计FastAPI本身是不关心数据操作的它只负责接收请求、调用你提供的函数、把返回值序列化成JSON。所以“数据操作”这部分的质量完全取决于你在FastAPI下面怎么组织数据库访问层。很多初学者习惯在路由函数里直接写SQLAlchemy查询项目小的时候没问题一旦路由变多、逻辑变复杂查询代码会散落在各个接口里之后想统一加缓存、统一处理事务、统一做分页都会非常痛苦。所以我建议在项目里把数据访问单独拆出一层大致目录结构是这样app/ ├── main.py # 应用入口包含路由注册与启动逻辑 ├── core/ │ ├── config.py # 环境变量与全局配置 │ └── database.py # 引擎、会话工厂、session依赖 ├── models/ # SQLAlchemy ORM模型 │ ├── __init__.py │ └── product.py ├── schemas/ # Pydantic模型用于接口入参与出参 │ ├── __init__.py │ └── product.py ├── crud/ # 数据操作层每个模型对应一个文件 │ ├── __init__.py │ └── product.py ├── api/ │ └── v1/ │ ├── __init__.py │ └── product.py # 路由层只负责HTTP相关处理 └── alembic/ # 数据库迁移脚本目录这种分层方式的好处是路由层不碰数据库只调用crud层函数crud层不依赖Request和Response对象只处理数据模型层不掺入业务逻辑。各层职责清晰后续上单元测试也会很容易——直接测crud层不需要启动HTTP服务。有一点必须说明这个结构不是FastAPI规定的是我从不同项目里总结出来的通用做法。如果项目特别小几个人临时合作写个脚本你完全可以简化但只要是预期能上线、会持续迭代的工程建议一开始就按这个来。中途拆分比一开始拆分痛苦得多。1.2 同步还是异步选型的关键判断依据FastAPI可以同时跑同步和异步接口但在数据库层面选型时需要先想清楚你的服务到底要面对什么样的并发场景。如果你的接口里只有非常简单的查询而且预估并发不高用同步SQLAlchemy加上线程池也能撑住但如果你做的是面向大量小请求的业务系统比如电商C端接口、消息推送服务那么事件循环被数据库阻塞一次整个进程能处理的请求数就会直线下降。我这边项目最终选择了异步SQLAlchemy具体来说是SQLAlchemy 2.0的异步特性配合asyncpg驱动连接PostgreSQL。异步写法和同步写法在模型定义上基本一致主要区别在查询执行时要用await session.execute(...)。实测下来在同样4核8G的容器里纯查询接口的QPS比同步方案高出将近一倍有IO等待的场景提升更明显。为什么选2.0版本而不是1.4因为2.0把传统的Query API真正淘汰了统一换成select()、update()、delete()这种2.0风格的语句。虽然1.4也有兼容模式但已经接触新项目的朋友没必要再学一套即将过时的写法。当然不是说所有项目都无脑上异步。如果你的业务大量涉及复杂报表、多表嵌套子查询、SQL调优异步带来的复杂度可能大于收益。这种场景下我反而建议用同步SQLAlchemy配合FastAPI的线程池或者更干脆一点直接用上文件里封装好的查询函数别让业务代码直接碰SQLAlchemy。把“选型”这个问题想明白了后面代码写起来会顺很多。2. 数据库连接与会话管理2.1 从配置文件到异步引擎所有配置都应该从环境变量读取不写死在代码里。用pydantic-settings管理配置是目前比较常见的做法它能自动读取.env文件并且在启动时做类型校验。下面这是config.py的快速实现from pydantic_settings import BaseSettings class Settings(BaseSettings): database_url: str postgresqlasyncpg://user:passlocalhost:5432/app_db echo_sql: bool False pool_size: int 10 max_overflow: int 20 pool_recycle: int 3600 class Config: env_file .env settings Settings()database_url里我用的是PostgreSQL的异步驱动asyncpg这也是生产环境最常用的组合。很多人会问为什么不用同步驱动psycopg2原因很简单asyncpg在异步场景下性能和协议支持都更好SQLAlchemy对它的支持也很成熟。接下来是database.py负责创建引擎和会话工厂from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession from sqlalchemy.orm import DeclarativeBase from core.config import settings engine create_async_engine( settings.database_url, echosettings.echo_sql, pool_sizesettings.pool_size, max_overflowsettings.max_overflow, pool_recyclesettings.pool_recycle, ) SessionLocal async_sessionmaker( bindengine, class_AsyncSession, expire_on_commitFalse, ) class Base(DeclarativeBase): pass这里有个容易被忽略的参数expire_on_commitFalse。默认情况下commit之后所有实例上的属性会失效下一次访问会重新触发SQL查询。在异步环境里这可能会在response序列化时意外触发懒加载然后抛MissingGreenlet异常。把它设为Falsecommit后对象属性仍然可读接口返回时会省很多麻烦。pool_size和max_overflow要与实际部署环境匹配。我给的是一个相对保守的起步配置总连接数约等于pool_size加max_overflow。这个值不能拍脑袋定我一般按“单个实例能够同时处理的最大并发数据库操作数”来估算。比如你的服务实例有4个worker每个worker并发跑20个请求那连接池理论上至少要80个连接否则请求会排队等待数据库连接。当然也不能无限调大数据库服务端有最大连接数限制连接太多反而会拖垮数据库。2.2 Session生命周期与FastAPI依赖注入Session的正确管理方式是每个请求创建一次请求结束关闭。不要全局复用一个session——这是新手最容易犯的错误。全局session在并发请求下会出现状态串线、事务相互干扰的问题。在FastAPI里做这件事非常顺手用依赖注入就行。我们在database.py里再定义一个生成器函数async def get_session() - AsyncSession: async with SessionLocal() as session: yield session然后在路由里这样使用from fastapi import APIRouter, Depends from sqlalchemy.ext.asyncio import AsyncSession from core.database import get_session from crud import product as product_crud router APIRouter() router.get(/products/{product_id}) async def get_product(product_id: int, db: AsyncSession Depends(get_session)): product await product_crud.get_product(db, product_id) return product这里yield之前的代码在创建sessionyield之后的代码在关闭session。async with块确保即使接口内部抛出异常session也会正常关闭连接归还给连接池。有一点要特别注意依赖注入里同一个session会在一个请求的整个生命周期中复用。这意味着如果你在接口里手动执行了session.commit()之后又继续操作数据库这些操作会开启新的事务和之前的事务并没有关联。因此我建议接口层面不要直接调commit而是让crud层里的写操作函数自己控制事务或者在crud层之外再包一层service逻辑统一提交。这块在后续事务部分会详细说。3. ORM模型设计与数据表映射3.1 类型选择和字段约束模型设计直接决定数据操作的复杂程度。用SQLAlchemy 2.0的Mapped写法定义一个商品表示例模型如下from datetime import datetime, timezone from typing import Optional from sqlalchemy import String, Integer, Numeric, DateTime, Text from sqlalchemy.orm import Mapped, mapped_column from core.database import Base class Product(Base): __tablename__ products id: Mapped[int] mapped_column(Integer, primary_keyTrue, autoincrementTrue) name: Mapped[str] mapped_column(String(128), nullableFalse, uniqueTrue, indexTrue) description: Mapped[Optional[str]] mapped_column(Text, defaultNone) price: Mapped[float] mapped_column(Numeric(10, 2), nullableFalse) stock: Mapped[int] mapped_column(Integer, nullableFalse, default0) created_at: Mapped[datetime] mapped_column(DateTime(timezoneTrue), server_defaultfunc.now()) updated_at: Mapped[datetime] mapped_column(DateTime(timezoneTrue), server_defaultfunc.now(), onupdatefunc.now())价格字段必须用Numeric而不是Float。Float在数据库里存储的是二进制浮点数做金额计算会出现莫名的小数误差Numeric则是精确数值类型。这在涉及金额的项目里是常识但还是会有人在这里踩坑。uniqueTrue和indexTrue可以组合使用。unique本身会创建索引单独设置index是为了让其他查询条件也能走索引。比如name字段加了unique约束按name精确查询已经能走索引了但如果经常按价格范围查、按创建时间排序这些字段就应该加index。created_at使用server_defaultfunc.now()让数据库来生成时间而不是我们在Python代码里手动赋值。好处是多个应用实例同时写入时时间由数据库统一生成不会因为机器时钟偏差导致排序混乱。updated_at用onupdatefunc.now()每次做UPDATE操作时数据库会自动刷新这个字段我们不需要在crud代码里维护它。3.2 Alembic迁移与表结构同步模型定义好之后需要把表结构同步到数据库。开发环境可以简单调用Base.metadata.create_all()但生产环境不建议这么做。原因是create_all只会创建缺失的表不会修改已经存在的表结构。项目迭代过程中要加字段、改类型create_all完全帮不上忙。迁移工具可以记录每一次结构变更并且支持升降级团队协作时其他人拉代码后也能快速把本地库迁移到最新的schema。我通常用Alembic配合异步引擎做迁移。初始化之后先改数据库连接配置指向async的URL再写迁移脚本。生成迁移的命令大致是alembic revision --autogenerate -m create products table alembic upgrade head自动生成的迁移文件会包含create_table、add_column之类的操作。我建议每次提交迁移前都把生成的脚本打开检查一遍因为autogenerate不是万能的索引名变更、约束调整、数据迁移这类操作经常需要手写补充。另外执行迁移时连接的是哪个库要在alembic配置里看仔细我吃过一次亏把测试环境的配置发到生产库上执行好在当时只是建了个空表没有造成数据损失。后续我把所有环境配置都改成了独立的.env文件并加上了迁移前的提示确认才彻底避免这类事故。4. 核心CRUD接口实操4.1 创建数据模型实例构造与错误处理写一条商品数据crud层里可以这样封装from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from models.product import Product from schemas.product import ProductCreate async def create_product(db: AsyncSession, data: ProductCreate) - Product: product Product( namedata.name, descriptiondata.description, pricedata.price, stockdata.stock, ) db.add(product) await db.commit() await db.refresh(product) return productcommit之后必须refresh一次。因为id、created_at、updated_at这些字段是由数据库生成的commit之后ORM实例上并没有这些值refresh会重新从数据库加载这些字段。如果不做refresh返回给前端的产品数据里id就是None前端往后拿这个id去请求详情就全部404了。这里的建议是create_product的errors处理放在接口层做。比如数据库里name有唯一约束插入重复名称会抛IntegrityError接口层可以捕获后返回409冲突。crud层尽量保持纯净不要混杂HTTP细节。4.2 查询数据主键、列表与分页单条查询很简单get和select两种写法都可以async def get_product(db: AsyncSession, product_id: int) - Optional[Product]: return await db.get(Product, product_id)列表加分页正确写法如下from sqlalchemy import select, func async def list_products( db: AsyncSession, page: int 1, page_size: int 20, keyword: str | None None, ): stmt select(Product) if keyword: stmt stmt.where(Product.name.ilike(f%{keyword}%)) total await db.scalar(select(func.count()).select_from(stmt.subquery())) stmt stmt.order_by(Product.id.desc()).offset((page - 1) * page_size).limit(page_size) rows (await db.execute(stmt)).scalars().all() return rows, totalcount子查询用了subquery而不是直接在原stmt上拼select(func.count())。因为原stmt可能包含order_by直接在它基础上加count会把排序也带进去虽然结果正确但会多一次无谓的排序开销。先包一层subquery再去掉排序逻辑效率会好一些。分页方式这里用offset/limit适用于中小规模数据。如果表数据量达到百万级更推荐keyset分页也叫cursor分页用上次返回的最后一条id作为位置标记。keyset分页不会随着页数增大而变慢避免了offset大页面时数据库全表扫描的问题。两者没有绝对好坏看数据量级和场景选择。4.3 更新与删除注意行级影响更新操作有两种方案一种是查出来再改属性一种是用update语句直接执行。第一种适合更新前需要校验业务规则的场景比如库存变更前要确认当前值第二种适合无条件的字段更新效率更高。同时更新多个字段时用update语句from sqlalchemy import update async def update_product( db: AsyncSession, product_id: int, data: ProductUpdate, ) - Optional[Product]: values data.model_dump(exclude_unsetTrue) if not values: return await get_product(db, product_id) stmt ( update(Product) .where(Product.id product_id) .values(**values) .returning(Product) ) result await db.execute(stmt) await db.commit() return result.scalar_one_or_none()exclude_unsetTrue表示只提交前端传了字段。比如前端只想改stock不传namename就不会被覆盖。这是PATCH语义的正确实现。删除操作同样可以用delete语句但要注意外键依赖。如果子表里有数据引用这个产品直接删除会抛外键约束错误。我通常会把删除做成软删除——加一个is_deleted字段查询时默认过滤掉已删除的数据。这样做的好处是数据可回溯也避免外键级联问题。到底用物理删除还是软删除取决于业务需求但作为一种数据操作方案我建议业务系统优先考虑软删除。5. 批量导入导出设计与事务处理5.1 商品批量导入从上传到入库的完整链路在实际项目中管理后台经常需要批量录入数据。与其一条条调用创建接口不如设计一个导入接口前端上传CSV或Excel文件后端解析、校验、批量写入。这个场景在电商后台尤其常见——供应商给的商品价格表、库存表动辄几千行逐条调用接口既不高效也容易把数据库和网络打满。导入接口的实现思路是接收上传文件保存为临时文件或直接读入内存。用pandas或openpyxl解析文件内容。逐行校验数据名称非空、价格大于0、库存非负、是否存在重复。将有效数据批量入库无效数据收集起来返回给前端方便用户下载错误报告。批量入库这一步值得展开。很多人会自然地写一个循环每条数据add一次最后commit。这个做法能用但性能很差。几千行的循环会导致几千次SQL flush耗时会放大很多倍。更合适的做法是用bulk操作from sqlalchemy.dialects.postgresql import insert as pg_insert async def bulk_upsert_products(db: AsyncSession, rows: list[dict]): stmt pg_insert(Product).values(rows) stmt stmt.on_conflict_do_update( index_elements[Product.name], set_{ price: stmt.excluded.price, stock: stmt.excluded.stock, description: stmt.excluded.description, }, ) await db.execute(stmt) await db.commit()这段代码用的是PostgreSQL的upsert语法按name判断是否已存在存在则更新价格、库存、描述不存在则插入。一条SQL解决几千行数据效率远高于循环逐条操作。有个性能经验值得分享批量入库时每批建议控制在500到1000行。太少起不到批量的效果太多则单条SQL过长可能触发数据库的SQL语句大小限制或参数数量限制。如果是几万行数据分批次执行不要一把梭。5.2 商品批量导出流式生成保护内存导出正好和导入相反数据量大时不要一次性把所有数据都加载到内存里。比如导出10万商品为CSV如果先全部查询出来再序列化应用进程内存可能直接翻几倍。正确做法是用流式响应分页读取数据逐行写入响应流。FastAPI可以用StreamingResponse实现这一点from fastapi.responses import StreamingResponse import csv, io async def export_products(db: AsyncSession, page_size: int 1000): def generate(): yield id,name,price,stock\n last_id 0 while True: stmt ( select(Product) .where(Product.id last_id) .order_by(Product.id) .limit(page_size) ) rows (await db.execute(stmt)).scalars().all() if not rows: break for p in rows: yield f{p.id},{p.name},{p.price},{p.stock}\n last_id rows[-1].id return StreamingResponse( generate(), media_typetext/csv, headers{Content-Disposition: attachment; filenameproducts.csv}, )这段代码用keyset分页方式替代offset避免深分页性能问题每次只取1000条处理完立即释放内存占用非常稳定。需要指出的是这里生成器函数里用了异步查询所以generate本身要能配合事件循环运行。实际项目中我一般会把查询逻辑抽到crud层让生成器函数只做拼装输出这样结构更清晰测试也更方便。5.3 事务边界与并发控制事务是数据操作里最核心的问题之一也是面试官最喜欢深挖的话题。一个事务要保证的是要么全部成功要么全部失败不能出现“库存扣了但订单没创建”这种中间状态。FastAPI里管理事务有两个层级。小范围用session.begin()块即可async with SessionLocal() as session: async with session.begin(): await session.execute(update_sql_1) await session.execute(update_sql_2)begin块结束时会自动commit块内任何一步抛异常都会自动rollback不需要手动编写commit/rollback逻辑。这是我在crud层写复杂写操作时最常用的方式。大范围跨多个crud函数的事务建议在service层组合。比如“下单”需要同时减库存和生成订单可以把两个crud函数放在同一个session里由service决定提交时机async def create_order(db: AsyncSession, user_id: int, items: list[dict]): async with db.begin(): order await order_crud.create_order(db, user_id) for item in items: await stock_crud.deduct_stock(db, item.product_id, item.quantity) await order_item_crud.create_order_item(db, order.id, item) return order这里db.begin()是外层事务内层crud函数里不能再有独立的commit。这是一个约定——写crud函数时保持“只写SQL不发言提交”由调用方控制事务边界。团队协作时最好把这条写进开发规范否则很容易出现两边都在控制事务导致保存点上下文错乱。并发控制这块最实用的手段是版本号乐观锁。在模型里加一个version字段version: Mapped[int] mapped_column(Integer, nullableFalse, default0)在SQLAlchemy的__mapper_args__中配置version_id_col指向它__mapper_args__ {version_id_col: version}这样每次UPDATE时SQLAlchemy会带上WHERE version 当前值更新成功后version加1。如果两个请求同时读到version1并尝试更新只有第一个能成功第二个update会影响到0行记录从而可以判断发生了并发冲突。这种方案很适合高并发扣库存、改配置这种场景。6. 常见问题与排查技巧实录6.1 MissingGreenlet错误异步懒加载的坑async环境里最常见的坑就是MissingGreenlet。触发原因通常是ORM实例在session关闭之后被访问了还未加载的关联属性或者懒加载触发时SQLAlchemy试图在隐式IO中执行查询但异步环境不允许这么做。它的典型报错长这样sqlalchemy.exc.MissingGreenlet: greenlet_spawn has not been called; cant call await_only() here. Was IO attempted in an unexpected place?发生场景一般是接口从数据库查到一条订单记录在session关闭后业务代码访问order.items这个关联属性。默认情况下items是懒加载的要等到访问时才去查数据库但此时session已经关了于是报错。解决办法是查询时主动声明加载关联关系用selectinloadfrom sqlalchemy.orm import selectinload stmt select(Order).options(selectinload(Order.items)).where(Order.id order_id)selectinload会用一条额外的IN查询把关联数据一次性加载避免懒加载。注意selectinload适合集合关联一对多、多对多单对象关联多对一用joinedload更合适。这个坑排查起来会稍微费时间因为报错不一定马上出现和请求的并发时序有关。我建议把“查询必带selectinload”写成crud层的代码规范保证所有会返回给前端的数据都提前加载好不依赖懒加载。6.2 连接池耗尽与超时接口突然大面积502或超时排查时发现数据库连接池被占满这是生产环境很棘手的问题之一。最典型的诱因session没有正确关闭。每个请求占用的连接都被高并发请求长时间持有数据库连接数被逐渐耗尽。排查思路是打印连接池使用状态或者直接查看数据库侧的活跃连接来源。修复方式是保证get_session依赖里用async with管理session确保请求结束、异常抛出时session都会关闭。第二个诱因是接口里的慢查询。某个查询没走索引执行时间几秒导致连接被长时间占用。排查方式是在SQLAlchemy配置里开启echo_sql或者更轻量的方式是记录慢查询日志定位后再给查询条件加合适的索引或改写查询逻辑。第三个诱因是连接池配置过小。如果代码没问题、慢查询也不多但并发就是上不去那就回到2.1节里的连接池配置把pool_size和max_overflow调大同时检查数据库服务端的max_connections。6.3 序列化递归与循环引用在写出参时Pydantic模型如果嵌套了父对象和子对象很容易写出来一个相互引用的结构导致序列化时无限递归最后栈溢出。解决方法是明确写出只向外序列化一层避免在Pydantic里把双向关系都暴露出来。比如Order模型里有itemsOrderItem模型里不要再嵌套Order只需要order_id。也就是说ORM模型之间可以双向关联但Pydantic模型必须设计成单向的树形结构。这个原则几乎所有序列化框架都适用。如果遇到ORM和Pydantic字段不完全匹配的情况比如ORM里是snake_case前端要camelCase可以写一个转换函数手动映射字段不要硬靠框架魔改。保持映射逻辑显式化代码反而更容易维护。6.4 查询N1问题N1问题的表现查询10条订单每条订单又各查一次关联表共执行了1次主查询加10次关联查询。接口响应时间随数据量线性恶化。排查方法是在SQLAlchemy的echo_sql日志里观察SQL条数。如果发现查询列表后跟着大量相同结构的SELECT那基本就是N1了。修复方法就是6.1节说的selectinload或joinedload在查询时一次性把关联数据加载出来。还有一个优化方向是查询时只select需要的列不select整个ORM对象。比如列表页只需要id、name、price就可以用select(Product.id, Product.name, Product.price)返回Row而不是完整ORM实例。这样数据库传输的数据量小ORM实例化和追踪的开销也少。但这个优化会让代码多写几行不是性能瓶颈的项目不需要过度优化。7. 数据操作的一些个人体会从我自己的项目经验来看数据操作层的核心不是把CRUD写出来而是把边界划清楚查询只管查询事务由调用方控制模型不掺业务逻辑Pydantic模型单向序列化。这四条规矩守住了无论项目怎么膨胀数据层都不会烂到不可维护。另外多说一句关于“抄作业”的建议网上有很多FastAPI项目的仓储代码GitHub上各种clean architecture模板也很多。我的建议是与其直接复制一个庞大的模板不如先理解template里每一层是做什么的然后根据自己项目的实际复杂度做取舍。手里有五六条路的时候别全铺上选最简单的、但能满足未来三个月扩展需求的那条就好。这套数据操作的方案在我这边的项目里已经跑了一年多上过生产、扛过促销峰值目前没有遇到新的绕不过去的坎。如果你按这个思路搭完数据层遇到什么奇怪的报错或者有其他更好的方案欢迎一起交流。
返回列表