第 11 章:分层架构与 CRUD 工程化
第 10 章把数据库接上了,但代码还是散的:一个
main.py里混着路由、SQL 和业务判断。本章给出全书统一的分层结构,把「HTTP 长什么样」「业务怎么算」「数据怎么存」三件事彻底分开。
学习目标
- 理解把逻辑全塞进路由函数会带来哪些具体后果
- 掌握本教程统一的目录结构与每层职责边界
- 能独立实现 Repository 层与 Service 层,并说清它们的差别
- 能判断事务边界该放在哪一层,并解释两种可行方案的取舍
- 能用
Annotated把get_db → repository → service依赖链组装起来 - 能实现通用的分页、排序、过滤,并输出一致的
Page[T]响应
11.1 为什么必须分层
先看一段「能跑」的代码:
@app.post("/users")
async def create_user(payload: UserCreate, session: AsyncSession = Depends(get_db)):
exists = (await session.execute(select(User).where(User.email == payload.email))).scalar_one_or_none()
if exists:
raise HTTPException(status_code=400, detail="邮箱已存在")
user = User(email=payload.email, hashed_password=hash_password(payload.password))
session.add(user)
await session.commit()
return user它没有语法错误,但它把四件事焊死在一个函数里。后果会随着项目规模逐个出现:
| 症状 | 具体表现 |
|---|---|
| 不可测试 | 想测「邮箱重复」这条规则,必须起一个 HTTP 客户端、造一份请求体、连一个数据库 |
| 不可复用 | 后台脚本要批量导入用户时,没法调用这段规则,只能把代码抄一份 |
| 事务边界混乱 | commit() 埋在路由里,跨多个资源的事务(下单 + 扣库存)无处安放 |
| 文件爆炸 | 每个资源 5 个端点 × 平均 30 行 = 一个路由文件上千行,IDE 跳转都卡 |
| 职责混淆 | HTTPException 把 HTTP 语义泄漏进业务层,同一段逻辑给 gRPC 或定时任务用时全部报错 |
分层的目标不是「看起来高级」,而是让每一层都能单独替换、单独测试。
11.2 统一目录结构
本教程全书使用下面这套结构,第 12 ~ 18 章的代码都往里放:
app/
├── main.py
├── core/ # 配置、安全、异常处理器
├── db/ # base.py、session.py
├── models/ # SQLAlchemy 模型
├── schemas/ # Pydantic 模型
├── repositories/ # 数据访问层
├── services/ # 业务逻辑层
└── api/
└── v1/ # 路由(APIRouter)依赖方向是单向的:api → services → repositories → db。反向依赖(例如 repository 导入 router)一律禁止,它会在运行时变成导入错误或启动即崩。
11.3 每层的职责边界
| 层 | 输入 | 输出 | 允许依赖 | 禁止做 |
|---|---|---|---|---|
路由层 api/v1 | HTTP 请求(路径/查询/请求体/请求头) | Pydantic schema | Service、Schemas、依赖函数 | 写 SQL、写业务规则、直接操作 session |
Service 层 services | schema 或领域对象、当前用户 | schema / 领域对象 / Page[T] | Repository、其他 Service、Schemas | 抛 HTTPException、读 Request、拼 URL |
Repository 层 repositories | 查询条件、模型实例 | ORM 模型 / 标量 / 元组 | AsyncSession、Models | 写业务规则、判断权限、提交事务 |
模型层 models | — | — | db/base | 依赖任何上层模块 |
Schema 层 schemas | — | — | 标准库、Pydantic | 导入 Service / Router |
一条实用判据:这一行代码如果换掉框架还要不要? 要 → 属于 Service 或 Repository;不要 → 留在路由层。
💡 Service 层不感知 HTTP,意味着它抛出的必须是领域异常(如
EmailAlreadyExistsError),由第 8 章写的异常处理器统一翻译成 409。这样同一段业务逻辑可以同时被 REST 接口、CLI 脚本和后台任务复用。
11.4 Repository:只做数据访问
# app/repositories/user.py
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.user import User
class UserRepository:
"""只负责读写数据,不做任何业务判断,不提交事务。"""
def __init__(self, session: AsyncSession) -> None:
self.session = session
async def get(self, user_id: int) -> User | None:
return await self.session.get(User, user_id) # 按主键取,命中身份映射时不发 SQL
async def get_by_email(self, email: str) -> User | None:
stmt = select(User).where(User.email == email).limit(1)
return (await self.session.execute(stmt)).scalar_one_or_none()
async def list(self, *, offset: int = 0, limit: int = 20, email_like: str | None = None) -> list[User]:
stmt = select(User)
if email_like:
stmt = stmt.where(User.email.ilike(f"%{email_like}%"))
stmt = stmt.order_by(User.id.desc()).offset(offset).limit(limit)
return list((await self.session.execute(stmt)).scalars().all())
async def count(self, *, email_like: str | None = None) -> int:
stmt = select(func.count()).select_from(User)
if email_like:
stmt = stmt.where(User.email.ilike(f"%{email_like}%"))
return (await self.session.execute(stmt)).scalar_one()
def add(self, user: User) -> User:
self.session.add(user)
return user
async def delete(self, user: User) -> None:
await self.session.delete(user)
async def flush(self) -> None:
await self.session.flush() # 只 flush,不 commit:提交权在 get_db(见 11.6)list / count 的过滤条件用同一套参数,是为了保证分页的 items 与 total 口径一致(见 11.7)。
11.5 Service:业务规则与编排
# app/services/user.py
from app.core.exceptions import EmailAlreadyExistsError, UserNotFoundError
from app.core.security import hash_password
from app.models.user import User
from app.repositories.user import UserRepository
from app.schemas.common import Page, PageParams
from app.schemas.user import UserCreate, UserUpdate
class UserService:
"""业务规则、编排多个 repository;只 flush 控制 SQL 时机,提交交给 get_db。"""
def __init__(self, repo: UserRepository) -> None:
self.repo = repo
async def create(self, data: UserCreate) -> User:
if await self.repo.get_by_email(data.email) is not None:
raise EmailAlreadyExistsError(data.email) # 领域异常,不是 HTTPException
user = User(email=data.email, hashed_password=hash_password(data.password))
self.repo.add(user)
await self.repo.flush() # 发 INSERT,拿到自增 id;提交由 get_db 统一完成
return user
async def get_or_404(self, user_id: int) -> User:
user = await self.repo.get(user_id)
if user is None:
raise UserNotFoundError(user_id)
return user
async def update(self, user_id: int, data: UserUpdate) -> User:
user = await self.get_or_404(user_id)
if data.email is not None and data.email != user.email:
if await self.repo.get_by_email(data.email) is not None:
raise EmailAlreadyExistsError(data.email)
user.email = data.email
for field, value in data.model_dump(exclude_unset=True).items():
setattr(user, field, value) # exclude_unset 保证「没传的字段」不被覆盖成 None
await self.repo.flush()
return user
async def delete(self, user_id: int) -> None:
user = await self.get_or_404(user_id)
await self.repo.delete(user)
await self.repo.flush() # 同样只 flush,等依赖退出时提交
async def list_all(self, limit: int = 100) -> list[User]:
"""管理端用:取一批用户。命名体现业务意图,而不是把 repo 暴露给路由。"""
return await self.repo.list(offset=0, limit=limit)三个值得注意的设计点:
- Service 不碰
session的查询 API,只用它flush()与rollback()。所有 SQL 都从 Repository 走,换存储实现时只需改一层。 model_dump(exclude_unset=True)是 PATCH 语义的关键。不加它,请求里没出现的字段会被当成None覆盖进数据库。- 异常是领域语言。
EmailAlreadyExistsError("a@b.com")比HTTPException(400, "邮箱已存在")携带的信息更准,映射方式由 API 层决定。
# app/core/exceptions.py
class AppError(Exception):
"""本应用所有领域异常的基类。"""
class UserNotFoundError(AppError):
def __init__(self, user_id: int) -> None:
self.user_id = user_id
super().__init__(f"用户 {user_id} 不存在")
class EmailAlreadyExistsError(AppError):
def __init__(self, email: str) -> None:
self.email = email
super().__init__(f"邮箱 {email} 已被注册")
class PermissionDeniedError(AppError):
"""已认证但权限不足,映射为 403。"""# app/api/errors.py(衔接第 8 章的异常处理器)
from fastapi import FastAPI, Request, status
from fastapi.responses import JSONResponse
from app.core.exceptions import AppError, EmailAlreadyExistsError, UserNotFoundError
def register_exception_handlers(app: FastAPI) -> None:
@app.exception_handler(UserNotFoundError)
async def _not_found(request: Request, exc: UserNotFoundError) -> JSONResponse:
return JSONResponse(
status_code=status.HTTP_404_NOT_FOUND,
content={"code": "USER_NOT_FOUND", "detail": str(exc)},
)
@app.exception_handler(EmailAlreadyExistsError)
async def _conflict(request: Request, exc: EmailAlreadyExistsError) -> JSONResponse:
# 409 Conflict:请求合法,但与现有资源状态冲突
return JSONResponse(
status_code=status.HTTP_409_CONFLICT,
content={"code": "EMAIL_EXISTS", "detail": str(exc)},
)💡 400 与 409 的界线:400 是「请求本身有问题」(缺字段、格式错、类型不对),409 是「请求没问题,但和当前数据状态冲突」(邮箱重复、版本号过期)。邮箱已存在属于后者,用 409 更准确。
11.6 事务边界:commit 放在哪一层
这是分层架构里争议最多的一个决定。先给结论:
commit()/rollback()只有一个归属——会话的所有者,也就是提供会话的get_db依赖。- Repository 只
flush(),绝不commit()。 - Service 同样只
flush(),用 flush 控制「SQL 什么时候发出去」,用异常控制「事务要不要作废」。
理由是业务规则的正确性往往跨越多次写操作:转账要先 repo.debit(from_id, amount) 再 repo.credit(to_id, amount),两次写必须同生共死。如果 commit() 下沉到 Repository 的每个方法里,扣款成功、入账失败时事务已经提交,钱就凭空消失了。Repository 无法知道「这一次写是不是一个完整业务动作的一部分」。
那么把 commit() 交给 Service 呢?它比 Repository 好,但比依赖统一提交差:一旦有多个 Service 方法参与同一次请求,谁先提交谁就把事务切断了,异常发生时留下半提交状态。把提交收敛到会话退出的那一刻,commit() 在整个请求里只有一个调用点。
# app/api/deps.py —— 会话的唯一所有者,也是唯一的事务边界
async def get_db() -> AsyncIterator[AsyncSession]:
async with SessionLocal() as session:
try:
yield session
await session.commit() # 正常返回 → 提交
except Exception:
await session.rollback() # 任何异常 → 回滚
raise两种可行方案的取舍:
| 方案 | 提交位置 | 优点 | 缺点 | 适用 |
|---|---|---|---|---|
| A:依赖统一提交(本书采用) | get_db 在依赖正常退出时 commit(),异常时 rollback() | 提交点全项目唯一,不可能出现半提交;Service 代码干净,不会「忘记提交」 | 事务边界隐式;Service 无法做「部分提交」;只读请求也带着一个事务 | 绝大多数 API 项目 |
| B:Service 显式提交 | Service 每个写方法末尾 await session.commit() | 事务范围肉眼可见,可以按业务需要划分多段事务 | 写方法一多就容易漏写;多 Service 协作时容易半提交 | 需要在一请求内分段提交的长流程 |
❌ 不推荐:
class UserRepository:
async def create(self, user: User) -> User:
self.session.add(user)
await self.session.commit() # 事务边界碎在数据访问层里
return user✅ 推荐:
class UserRepository:
async def create(self, user: User) -> User:
self.session.add(user)
await self.session.flush() # 只把 SQL 发出去,提交由会话所有者决定
return user
class UserService:
async def register(self, data: UserCreate) -> User:
# 不 commit:外层 get_db 会在请求成功时统一提交
return await self.repo.create(User(email=data.email, hashed_password=hash_password(data.password)))方案 A 下 Service 依然需要 flush(),因为业务规则常常依赖数据库的反馈:
flush()之后自增主键与数据库默认值才可用,Service 才能返回带id的对象。- 唯一约束、外键约束的违反发生在
flush()时(而不是commit()),Service 才能就近把IntegrityError翻译成领域异常。
无论选哪种方案,都要遵守一条铁律:一次 HTTP 请求只用一个 AsyncSession,绝不能出现两个会话各自提交。两个会话 = 两个事务 = 没有原子性。
11.7 依赖链组装
# app/api/deps.py
from collections.abc import AsyncIterator
from typing import Annotated
from fastapi import Depends
from sqlalchemy.ext.asyncio import AsyncSession
from app.db.session import SessionLocal
from app.repositories.user import UserRepository
from app.services.user import UserService
async def get_db() -> AsyncIterator[AsyncSession]: # 第 10 章的会话依赖
async with SessionLocal() as session:
try:
yield session
await session.commit() # 一个请求 = 一个事务,唯一的提交点
except Exception:
await session.rollback()
raise
DbSession = Annotated[AsyncSession, Depends(get_db)]
def get_user_repository(session: DbSession) -> UserRepository:
return UserRepository(session)
def get_user_service(repo: Annotated[UserRepository, Depends(get_user_repository)]) -> UserService:
return UserService(repo) # 提交由 get_db 负责,Service 无需持有 session
UserServiceDep = Annotated[UserService, Depends(get_user_service)]依赖链上的每一环都是独立可替换的接缝:第 16 章测试时,想换数据库就 app.dependency_overrides[get_db],想绕过业务规则就覆盖 get_user_service,想造假数据就覆盖 get_user_repository。这就是依赖注入在工程化上的全部意义。
11.8 通用分页、排序与过滤
分页参数几乎每个列表接口都要用,抽成可复用依赖:
# app/schemas/common.py
from typing import Generic, TypeVar
from pydantic import BaseModel, Field
T = TypeVar("T")
class PageParams(BaseModel):
page: int = Field(default=1, ge=1, description="从 1 开始的页码")
size: int = Field(default=20, ge=1, le=100, description="每页条数")
@property
def offset(self) -> int:
return (self.page - 1) * self.size
class Page(BaseModel, Generic[T]):
items: list[T]
total: int
page: int
size: int
pages: int # 总页数,前端做分页器直接可用# app/api/deps.py(续)
from fastapi import Query
from app.schemas.common import PageParams
def get_page_params(
page: Annotated[int, Query(ge=1)] = 1,
size: Annotated[int, Query(ge=1, le=100)] = 20,
) -> PageParams:
return PageParams(page=page, size=size)
PageParamsDep = Annotated[PageParams, Depends(get_page_params)]Service 侧把「列表 + 总数」组装成 Page:
# app/services/user.py(续)
# app/services/user.py(续)
async def list_paginated(self, params: PageParams, email_like: str | None = None) -> Page[User]:
items = await self.repo.list(offset=params.offset, limit=params.size, email_like=email_like)
total = await self.repo.count(email_like=email_like)
pages = (total + params.size - 1) // params.size if total else 0 # 向上取整
return Page[User](items=items, total=total, page=params.page, size=params.size, pages=pages)路由层用泛型响应模型,OpenAPI 会生成正确的 schema:
@router.get("", response_model=Page[UserOut])
async def list_users(service: UserServiceDep, params: PageParamsDep) -> Page[UserOut]:
page = await service.list_paginated(params)
# Page 是泛型模型,Session 执行完这里不再有惰性 IO,安全转成出参 schema
return Page[UserOut](
items=[UserOut.model_validate(item) for item in page.items],
total=page.total,
page=page.page,
size=page.size,
pages=page.pages,
)排序字段要白名单,绝不能让用户传的字符串直接进 order_by:
❌ 不推荐:
stmt = select(User).order_by(text(f"{sort_by} {order}")) # 注入风险 + 数据库报错泄漏✅ 推荐:
SORTABLE = {"id": User.id, "email": User.email, "created_at": User.created_at}
def apply_sort(stmt, sort_by: str, order: str):
column = SORTABLE.get(sort_by, User.id) # 未知字段回落到默认排序
return stmt.order_by(column.desc() if order == "desc" else column.asc())⚠️
total与items是两条查询,中间可能有并发写入,导致total与items长度不自洽(第 3 页只剩 2 条)。对一致性要求高的场景,用窗口函数count(*) OVER ()在一条 SQL 里同时取回两者。
11.9 端到端 users 模块
① schemas
# app/schemas/user.py
from datetime import datetime
from pydantic import BaseModel, ConfigDict, EmailStr, Field
class UserBase(BaseModel):
email: EmailStr
class UserCreate(UserBase):
password: str = Field(min_length=8, max_length=128)
class UserUpdate(BaseModel):
email: EmailStr | None = None
nickname: str | None = Field(default=None, max_length=50)
class UserOut(UserBase):
model_config = ConfigDict(from_attributes=True) # 允许直接从 ORM 对象构造
id: int
nickname: str | None = None
is_active: bool
created_at: datetime② models:app/models/user.py 与 post.py 见第 10 章,User 需补一个 nickname: Mapped[str | None]。
③ repositories:见 11.4。
④ services:见 11.5。
⑤ router
# app/api/v1/users.py
from typing import Annotated
from fastapi import APIRouter, Depends, Query, status
from app.api.deps import PageParamsDep, UserServiceDep
from app.schemas.common import Page
from app.schemas.user import UserCreate, UserOut, UserUpdate
router = APIRouter(prefix="/users", tags=["users"])
@router.get("", response_model=Page[UserOut], summary="分页查询用户")
async def list_users(
service: UserServiceDep,
params: PageParamsDep,
email_like: Annotated[str | None, Query(max_length=100)] = None,
) -> Page:
return await service.list_paginated(params, email_like)
@router.post("", response_model=UserOut, status_code=status.HTTP_201_CREATED, summary="注册用户")
async def create_user(payload: UserCreate, service: UserServiceDep) -> UserOut:
user = await service.create(payload)
return UserOut.model_validate(user)
@router.get("/{user_id}", response_model=UserOut, summary="查询单个用户")
async def get_user(user_id: int, service: UserServiceDep) -> UserOut:
return UserOut.model_validate(await service.get_or_404(user_id))
@router.patch("/{user_id}", response_model=UserOut, summary="局部更新用户")
async def update_user(user_id: int, payload: UserUpdate, service: UserServiceDep) -> UserOut:
return UserOut.model_validate(await service.update(user_id, payload))
@router.delete("/{user_id}", status_code=status.HTTP_204_NO_CONTENT, summary="删除用户")
async def delete_user(user_id: int, service: UserServiceDep) -> None:
await service.delete(user_id)⑥ 聚合路由与挂载
# app/api/v1/__init__.py
from fastapi import APIRouter
from app.api.v1 import users
api_router = APIRouter(prefix="/api/v1")
api_router.include_router(users.router)# app/main.py
from contextlib import asynccontextmanager
from fastapi import FastAPI
from app.api.errors import register_exception_handlers
from app.api.v1 import api_router
from app.core.config import settings
from app.db.base import Base
from app.db.session import engine
@asynccontextmanager
async def lifespan(app: FastAPI):
if settings.ENVIRONMENT == "dev":
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
yield
await engine.dispose()
app = FastAPI(title=settings.PROJECT_NAME, lifespan=lifespan)
register_exception_handlers(app)
app.include_router(api_router)main.py 只做三件事:装配 lifespan、注册异常处理器、挂载路由。任何业务代码都不应该出现在这里——它变长了,通常意味着某段逻辑放错了层。
常见坑与排查
| 现象 | 原因 | 解决 |
|---|---|---|
启动报 ImportError: cannot import name 'UserOut' from partially initialized module | 循环导入:schemas 导入了 models,models 又间接导入 schemas | 依赖只能单向:api → services → repositories → models;类型注解用字符串或 TYPE_CHECKING 延迟导入 |
| Service 里的规则无法被脚本/任务复用 | Service 抛了 HTTPException,把 HTTP 语义焊进业务层 | 改抛领域异常,在 app/api/errors.py 里统一映射成状态码 |
Repository 里每个方法都只有一行 session.execute | 「薄封装」没有附加价值,只是多一层跳转 | 让它承担明确职责:统一过滤条件、只暴露必要的查询、隔离 ORM 细节(如 selectinload) |
报错栈里全是 deps.py,定位不到业务代码 | 依赖链过深(get_db → repo → service → sub-service → ...) | 依赖层级不超过 3 层;同一层内的协作靠构造参数传递而非再套 Depends |
列表接口 total 与实际条数对不上 | total 与 items 是两条 SQL,中间数据变了;或过滤条件口径不一致 | 让 list 与 count 共用同一个过滤函数;强一致场景用 count(*) OVER () 一次查回 |
PATCH 之后没传的字段被清空 | 用了 model_dump(),未提供的字段以 None 参与赋值 | 用 model_dump(exclude_unset=True);再配 exclude_none=True 视语义而定 |
| 接口偶发「邮箱已存在」但数据库查不到 | 唯一约束靠「先查后写」,并发下两个请求同时通过检查 | 数据库层加 unique=True(已加)+ 捕获 IntegrityError 转成 409,检查只是友好提示 |
路由文件越来越长、main.py 也膨胀 | 没有按资源拆分 router,聚合逻辑与业务逻辑混写 | 每个资源一个 api/v1/<resource>.py,api/v1/__init__.py 只做 include_router |
同一请求里出现多个 commit(),异常时留下半提交 | 提交点散落在多个 Service 方法里 | 把 commit() / rollback() 收敛到 get_db 一个地方,Service 只 flush() |
| 请求报错后数据「写进去了一半」 | 中途某处已 commit(),后续步骤失败却无法回滚 | 同上:整请求一个事务,异常统一 rollback();确需分段提交的长流程才用 Service 显式提交 |
本章小结
| 要点 | 说明 |
|---|---|
| 分层动因 | 让每层可单独替换、单独测试;避免路由文件膨胀与事务边界碎裂 |
| 目录结构 | core / db / models / schemas / repositories / services / api/v1 |
| 依赖方向 | 严格单向 api → services → repositories → db,反向依赖必成循环导入 |
| 路由层 | 解析请求、调用 Service、返回 schema;不写 SQL、不做业务判断 |
| Service 层 | 业务规则、编排 Repository;不感知 HTTP,抛领域异常,只 flush() 不 commit() |
| Repository 层 | 纯数据访问,接收/返回模型;只 flush(),绝不 commit() |
| 事务边界 | get_db 持有唯一的事务边界:正常退出 commit(),异常 rollback();Service 显式提交是需要分段事务时的备选方案 |
| 依赖组装 | get_db → get_user_repository → get_user_service,全部用 Annotated 别名,测试可逐环覆盖 |
| 分页 | PageParams 依赖 + 泛型 Page[T];list 与 count 必须共用过滤条件 |
| 排序 | 字段白名单映射到列对象,禁止把用户输入拼进 order_by |
练习题
为
Post资源补一整套分层代码(schemas / repository / service / router),要求列表接口支持按author_id过滤与按created_at排序,并挂到/api/v1/posts。实现
PageParams的「游标分页」变体:用created_at + id作为游标,返回next_cursor。说明它相比页码分页在什么场景下更合适,以及失去「跳转到第 N 页」能力的原因。用
app.dependency_overrides写一个测试,把get_user_service换成返回固定列表的桩对象,验证GET /api/v1/users的响应结构与pages字段计算正确。「先查邮箱是否存在,再插入」在并发下会失效。请写出用
IntegrityError兜底的完整 Service 方法,并说明为什么数据库唯一约束不可省略。下面这段 Service 代码有三处分层问题,请指出并改写:
class PostService:
async def publish(self, post_id: int, session: AsyncSession) -> dict:
post = await session.get(Post, post_id)
if post is None:
raise HTTPException(status_code=404, detail="文章不存在")
post.published = True
await session.commit()
return {"id": post.id, "published": post.published}下一章预告
分层结构让 CRUD 有了骨架,但现在的接口谁都能调。下一章补上最后一块骨头:认证与授权——密码怎么存、JWT 怎么签发校验、权限怎么按角色收口。