17 数据操作之新增,三步写进数据库
前置阅读: 16-数据操作之查询
关键词:db.add,await db.commit,await db.refresh, 事务
难度: ★★★☆☆
场景导入:一条 POST /books,背后发生了什么?
第 16 章把查询从内存搬到了 ORM 异步栈。新增紧随其后——语义同样是"把 Python 对象变成数据库新行",但它比查询多了一个关键步骤:事务边界。你必须保证 INSERT 落库、自增主键回填、默认值同步这三件事在同一个事务里完成。
本期要把第 06 章占位的 POST /books 改造为基于 ORM 的真实端点,在 app/dao/book_dao.py 中实现 create_book。核心是三步:db.add(book) 把对象登记到 Session→ await db.commit() 触发 flush + 提交事务 → await db.refresh(book) 拉回数据库生成的字段。每一步都有明确的 await 语义和异常路径,搞清楚了新增,更新和删除就是换一套 SQL 的事。
原理解析:add 是同步的,commit 和 refresh 必须 await
新增的三步操作在异步栈下有重要的语义差异:
db.add 是 SQLAlchemy Core 的纯 Python API——它只操作内存中的 Session.new 集合,不碰网络。在异步栈下也保持同步语义,加 await 反而报错。
await db.commit() 才是真正的 I/O 边界。AsyncSession.commit 是协程,内部先执行 flush(把 Session.new 和 Session.dirty 里的改动翻译为 SQL),再执行底层事务的 COMMIT。如果 flush 阶段抛出 IntegrityError(比如 ISBN 重复键 1062、外键约束失败 1452),commit 不会执行,事务保持打开——需要业务代码显式 await db.rollback() 才能让连接恢复可用。
await db.refresh(book) 重新从数据库加载该对象的所有字段,主要用于同步数据库端生成的默认值。MySQL 的自增主键在 INSERT 返回时就已回填,但 server_default(如 created_at)、触发器生成的值则需要再查一次才能拿到。在 expire_on_commit=False 的设置下,未 refresh 的属性保持 commit 前的值——“按需同步"而非"每次 commit 触发隐式 SELECT”。
整个新增流程的时序如下:
(图注:db.add 只修改 Session 内部集合,await db.commit() 触发 flush + 提交事务,await db.refresh(book) 拉回数据库生成的字段。)
事务的异常路径同样重要。db.add 之后、await db.commit() 之前如果发生任何异常,get_async_db 的 async with 退出阶段会自动 rollback,连接安全归还连接池:
(图注:commit 是事务出口,refresh 是按需拉取数据库生成字段,rollback 是异常路径上的清理动作。)
async with session.begin() 和显式 await db.commit 是两种事务控制风格。前者把事务边界写在上下文里,异常路径自动 rollback;后者更显式,commit 写在 DAO 函数内部,事务粒度与业务逻辑对齐。本系列采用显式 await db.commit,因为 DAO 函数可能包含多条语句,事务边界由开发者控制更清晰。
代码实现:POST /books 的真实实现
DAO 层新增 create_book,只接收业务字段,status 由服务端默认设置:
# app/dao/book_dao.py(增量追加)
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.book import Book, BookStatus
async def create_book(
db: AsyncSession,
*,
title: str,
author_id: int,
isbn: str,
) -> Book:
"""新增图书,默认状态为 AVAILABLE。"""
book = Book(
title=title,
author_id=author_id,
isbn=isbn,
status=BookStatus.AVAILABLE,
)
db.add(book)
await db.commit()
await db.refresh(book)
return book
Schema 层定义 BookCreate,字段约束与第 14 章的模型定义对齐:
# app/schemas/book.py(节选)
from pydantic import BaseModel, Field
from app.models.book import BookStatus
class BookCreate(BaseModel):
title: str = Field(min_length=1, max_length=128)
author_id: int = Field(gt=0)
isbn: str = Field(min_length=13, max_length=13)
min_length=13, max_length=13 与模型层的 String(13) 对齐——Pydantic 在请求体校验阶段就拦截不合法的长度,FastAPI 返回 422,不会把脏数据传到数据库层触发 DataError。这是"校验前置"的典型实践。
路由层追加 POST /books handler,返回 201 Created:
# app/api/books.py(追加路由)
from app.schemas.book import BookCreate
@router.post("", status_code=201)
async def create_book(payload: BookCreate, db: DBDep) -> Book:
"""新增图书。"""
return await book_dao.create_book(db, **payload.model_dump())
payload.model_dump() 把 Pydantic 模型展平为关键字参数字典,DAO 函数签名保持扁平——不依赖 Schema 类型,降低耦合。handler 是 async def,因为它需要 await DAO 函数。返回的 Book ORM 实例由 response_model=BookOut 自动序列化(第 06 章配置的 from_attributes=True)。
避坑指南
await db.refresh(book)按需使用。如果模型没有依赖任何数据库端默认值(server_default、触发器),可以省略这步,减少一次 SELECT。但如果有created_at等数据库端生成的字段,必须显式 refresh。跨请求不要复用 ORM 对象。
AsyncSession关闭后,book 实例处于 detached 状态,访问属性可能触发DetachedInstanceError或读到陈旧值。跨请求只传book.id,下次请求重新await db.get(Book, book_id)。显式 commit vs 上下文管理器。
async with session.begin():更简洁但适合单条语句场景;显式await db.commit()在包含多条 DAO 调用的事务中让边界更清晰。本系列后续的借书事务(改 Book + 插 Borrow)将演示同一 commit 内多表写入。IntegrityError必须 rollback。虽然get_async_db的async with退出阶段会兜底回滚,但捕获IntegrityError后显式await db.rollback()让意图更清晰、调试更容易。
面试 QA
Q1 [原理]: SQLAlchemy 中 flush 与 commit 的差异?异步栈下都需要 await 吗?
flush 把 Session 累积的改动翻译为 SQL 并发送到数据库,但不提交事务——其他连接看不到改动。commit 在 flush 之后提交事务,改动持久化并释放锁。两者在异步栈下都是协程,都必须 await,但只有 commit 是高频用法——flush 通常被 commit 隐式调用,业务代码很少显式 await db.flush()。
只有一种场景需要显式 flush:需要拿到自增主键但不想立即提交(比如批量导入时先 flush 拿 id,最后统一 commit)。本系列第 17 章没有用到,因为单条 INSERT 走 await db.commit() 即可。
Q2 [项目]: FastAPI 新增接口的事务边界如何把控?异常路径如何处理?
事务边界由 get_async_db 的 async with 和 DAO 函数内的 await db.commit() 共同决定。请求进入时框架创建 Session,DAO 在函数体内显式 commit;响应序列化后框架回到 async with 退出阶段关闭 Session。异常路径上,即使 DAO 不显式 await db.rollback(),async with 退出阶段也会自动回滚——但显式 rollback 让代码意图更清晰。
典型异常处理形态:
async def create_book_safe(db: AsyncSession, **kwargs) -> Book:
book = Book(**kwargs)
db.add(book)
try:
await db.commit()
except IntegrityError:
await db.rollback()
raise
await db.refresh(book)
return book
第 08 章的全局异常处理器负责把 IntegrityError 映射为 HTTP 状态码(1062 → 409,1452 → 422),DAO 层只需要 raise。
小结
本期把 POST /books 改造为基于 ORM 的异步新增,核心三步:db.add(book)(同步,登记到 Session)、await db.commit()(异步,flush + 提交事务)、await db.refresh(book)(异步,拉回数据库生成字段)。事务边界由 get_async_db 的 async with 和 DAO 的显式 await commit 共同把控,异常路径由全局 IntegrityError 处理器统一收口。
下一篇《18 数据操作之更新》将进入更新场景——讨论 update(Book).where(...).values(...) 的 Core 层写法,以及 version 列 + rowcount 校验实现的乐观锁。