第 10 章:数据库集成(SQLAlchemy 2.0)
前九章的应用都是「无状态」的:请求进来、算完、返回、忘掉。本章解决持久化——用 SQLAlchemy 2.0 异步模式接上数据库,把引擎、会话、模型、迁移一次讲清。
学习目标
- 理解异步驱动与同步驱动的差别,能按环境写出正确的连接串
- 掌握 2.0 声明式写法(
DeclarativeBase/Mapped/mapped_column),能识别 1.x 旧写法 - 能用
async_sessionmaker+yield依赖管理会话生命周期(衔接第 6 章) - 掌握异步 CRUD 动作:
select()/execute()/scalars()/flush()/refresh() - 理解异步下懒加载为何必然失败,并会用
selectinload预加载关系 - 能用 Alembic 生成并执行迁移,清楚
create_all与迁移的边界
10.1 选型与依赖
FastAPI 端点默认跑在事件循环里,一个 worker 进程用一个线程处理成百上千连接。在这里做同步数据库 I/O,整个进程都会被那次 I/O 卡住,异步的并发能力就白给了。
| 组合 | 驱动示例 | 连接串 | 适用场景 |
|---|---|---|---|
| 同步 ORM | psycopg2 / sqlite3 | postgresql://... | 传统 WSGI 应用、Celery 任务 |
同步 ORM + def 端点 | 同上 | 同上 | 过渡期,FastAPI 会把 def 丢进线程池 |
| 异步 ORM | asyncpg / aiosqlite | postgresql+asyncpg://... | 本教程统一采用 |
选它的理由:异步会话能直接用在 async def 端点、WebSocket、后台任务、httpx 调用之后,不必来回切线程池;AsyncSession / async_sessionmaker / create_async_engine 都是官方原生 API;连接池行为完全显式,没有隐式线程抢占。
| 环境 | 数据库 | 驱动包 | 连接串 |
|---|---|---|---|
| 开发 | SQLite 文件 | aiosqlite | sqlite+aiosqlite:///./app.db |
| 测试 | SQLite 内存 | aiosqlite | sqlite+aiosqlite:///:memory: |
| 生产 | PostgreSQL | asyncpg | postgresql+asyncpg://user:pwd@host:5432/db |
# 基础三件套:异步扩展 + 开发用 SQLite 驱动 + 迁移工具
uv add "sqlalchemy[asyncio]" aiosqlite alembic
# 生产环境再加 PostgreSQL 异步驱动
uv add asyncpgsqlalchemy[asyncio] 的 extra 会带上 greenlet。它是异步 ORM 的硬依赖,缺它不会在导入时报错,而是在第一次查库时抛 MissingGreenlet——典型「装了但装漏了」的坑。
连接串格式为 方言+驱动://用户:密码@主机:端口/库名:
sqlite+aiosqlite:///./app.db # 三个斜杠:当前目录下的相对路径
sqlite+aiosqlite:////srv/data/app.db # 四个斜杠:绝对路径 /srv/data/app.db
postgresql+asyncpg://postgres:secret@127.0.0.1:5432/blog⚠️ 连接串不写驱动名是最常见的错误。
postgresql://...会让 SQLAlchemy 去找同步驱动psycopg2,然后在await处报错。必须写postgresql+asyncpg://。
10.2 创建异步引擎
# app/db/session.py
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
DATABASE_URL = "sqlite+aiosqlite:///./app.db"
engine = create_async_engine(
DATABASE_URL,
echo=True, # 打印 SQL,开发期打开,生产务必关掉
pool_pre_ping=True, # 取连接前先探活,剔除被数据库单方面断开的死连接
pool_size=10, # 池中常驻连接数
max_overflow=20, # 峰值时额外临时创建的连接数上限
pool_recycle=1800, # 连接存活超过 1800 秒就重建,规避服务端空闲超时
)| 参数 | 作用 | 建议值 |
|---|---|---|
echo | 把 SQL 打到日志 | 开发 True,生产 False |
pool_size | 常驻连接数 | 与数据库 max_connections、worker 数一起算 |
max_overflow | 溢出连接上限 | 默认 10;可调大,但别超过数据库承受能力 |
pool_pre_ping | 取连接前探活 | 生产必开,否则会拿到被防火墙掐断的死连接 |
pool_recycle | 连接最大存活秒数 | 小于网关空闲超时(常见 1800) |
pool_timeout | 等空闲连接的超时 | 默认 30 秒,超时抛 TimeoutError |
SQLite 是另一套规则:文件库默认用 NullPool(每次新建、用完关闭),此时 pool_size 不生效,强行指定会直接报错。测试用共享内存库需要 StaticPool:
from sqlalchemy.pool import StaticPool
test_engine = create_async_engine(
"sqlite+aiosqlite:///:memory:",
connect_args={"check_same_thread": False},
poolclass=StaticPool, # 共享单连接,否则每个连接看到的是一个空库
)10.3 模型:2.0 声明式写法
2.0 用 DeclarativeBase 取代 1.x 的 declarative_base(),用 Mapped[...] + mapped_column() 取代裸 Column(...):
# app/db/base.py
from datetime import datetime
from sqlalchemy import DateTime, func
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
class Base(DeclarativeBase):
"""全项目共用的声明式基类,Alembic 的 target_metadata 指向它。"""
class TimestampMixin:
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
)❌ 不推荐(1.x 旧写法):
from sqlalchemy import Column, Integer, String
from sqlalchemy.orm import declarative_base
Base = declarative_base()
class User(Base):
__tablename__ = "users"
id = Column(Integer, primary_key=True, index=True)
nickname = Column(String, nullable=True) # 可空性埋在参数里,类型系统看不见✅ 推荐(2.0 写法):
from sqlalchemy import String
from sqlalchemy.orm import Mapped, mapped_column
class User(Base):
__tablename__ = "users"
id: Mapped[int] = mapped_column(primary_key=True, index=True)
nickname: Mapped[str | None] = mapped_column(String(50), default=None)差别不只是语法糖:Mapped[str | None] 让「这一列可为 NULL」成为类型系统的一部分,IDE 与 mypy 都能查出来;Mapped[int] 等价于 nullable=False。注解既是 Python 类型,也是 DDL 依据。金额用 Mapped[Decimal] = mapped_column(Numeric(10, 2)),长文本用 Text,时间用 DateTime(timezone=True)。
⚠️ 金额永远用
Numeric/Decimal,不要用float。浮点误差在对账时是灾难。
10.4 一对多关系:User 与 Post
# app/models/user.py
from sqlalchemy import JSON, String
from sqlalchemy.orm import Mapped, mapped_column, relationship
from app.db.base import Base, TimestampMixin
class User(TimestampMixin, Base):
__tablename__ = "users"
id: Mapped[int] = mapped_column(primary_key=True, index=True)
email: Mapped[str] = mapped_column(String(255), unique=True, index=True)
hashed_password: Mapped[str] = mapped_column(String(255))
nickname: Mapped[str | None] = mapped_column(String(50), default=None, nullable=True)
is_active: Mapped[bool] = mapped_column(default=True)
roles: Mapped[list[str]] = mapped_column(JSON, default=list) # 第 12 章的权限依赖会用到
# 只创建 Post 表,避免与 post.py 循环导入;Post 通过 back_populates 反向指回来
posts: Mapped[list["Post"]] = relationship(
back_populates="author",
cascade="all, delete-orphan",
passive_deletes=True,
)# app/models/post.py
from sqlalchemy import ForeignKey, String, Text
from sqlalchemy.orm import Mapped, mapped_column, relationship
from app.db.base import Base, TimestampMixin
class Post(TimestampMixin, Base):
__tablename__ = "posts"
id: Mapped[int] = mapped_column(primary_key=True, index=True)
title: Mapped[str] = mapped_column(String(200), index=True)
content: Mapped[str] = mapped_column(Text)
author_id: Mapped[int] = mapped_column(ForeignKey("users.id", ondelete="CASCADE"), index=True)
author: Mapped["User"] = relationship(back_populates="posts")三处细节:
back_populates必须双向写全。只写一边,另一边不会同步,user.posts与post.author会指向不同的 Python 对象。cascade="all, delete-orphan"是 ORM 层行为:删除user时,会话自动对已加载的post发DELETE;把post从user.posts移除,它也会被删除。ondelete="CASCADE"是数据库层行为,写在 DDL 上。两者不是替代关系,最好都写:ORM 路径干净,绕过 ORM 的直接 SQL 删除也不会留下孤儿行。
⚠️ SQLite 默认不强制外键约束。要让
ondelete生效,得为每条连接打开 pragma:在@event.listens_for(engine.sync_engine, "connect")里执行PRAGMA foreign_keys=ON。
10.5 会话工厂与 expire_on_commit=False
# app/db/session.py(续)
SessionLocal = async_sessionmaker(
bind=engine,
class_=AsyncSession,
expire_on_commit=False, # 异步下几乎是必须的
autoflush=False, # 关掉隐式 flush,让 flush 时机可控
)expire_on_commit 默认是 True:commit() 之后把会话里所有对象标记为过期,下次访问属性时重新发一条 SELECT。同步世界阻塞一下就好;但在 async def 里属性访问无法 await,SQLAlchemy 手上只有异步驱动,于是直接抛错:
sqlalchemy.exc.MissingGreenlet: greenlet_spawn has not been called;
can't call await_only() here. Was IO attempted in an unexpected place?❌ 不推荐:
async def create_user(session: AsyncSession, email: str) -> User:
user = User(email=email)
session.add(user)
await session.commit()
print(user.id) # commit 后对象已过期,这里触发重载 → MissingGreenlet
return user # 路由层序列化时同样炸✅ 推荐:
SessionLocal = async_sessionmaker(bind=engine, expire_on_commit=False)
async def create_user(session: AsyncSession, email: str) -> User:
user = User(email=email)
session.add(user)
await session.commit()
await session.refresh(user) # 显式取回数据库生成的默认值与自增 id
return user # 属性仍可安全访问(这里保留 commit() 是为了演示「提交后如何安全取回字段」;在分层项目里请把提交交给 get_db 依赖,见 10.6 与第 11 章。)
expire_on_commit=False 不是「放弃一致性」,只是不在提交后主动作废对象。真需要数据库最新值时用 await session.refresh(obj) 显式取回,意图反而更清楚;只刷新部分列可传 attribute_names=[...]。
10.6 会话依赖注入(衔接第 6 章)
第 6 章讲过 yield 依赖:yield 之前准备资源,yield 之后清理。数据库会话是它最经典的用例。
# app/db/session.py(续)
from collections.abc import AsyncIterator
async def get_db() -> AsyncIterator[AsyncSession]:
async with SessionLocal() as session:
try:
yield session
await session.commit() # 路由正常返回 → 提交本次请求的全部改动
except Exception:
await session.rollback() # 任何异常 → 整体回滚
raiseasync with SessionLocal()已负责关闭,不要手写session.close()。commit()只在这里出现一次:整个请求是一个事务,yield之后的代码在路由返回后执行,成功则提交、出错则回滚。这样就不会出现「前半段写进去了、后半段失败」的半提交。这一决策的完整取舍见第 11 章 11.6。rollback()是必需的兜底,否则异常路径会把一个未结束的事务留在连接上,连接归还连接池时状态是脏的。- 绝不要把
session存成模块级全局变量复用。AsyncSession既不线程安全也不协程安全,跨请求共享会导致InterfaceError与数据串味。
10.7 异步 CRUD 标准动作
# app/repositories/user_repository.py(第 11 章会正式分层,这里先看动作)
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.user import User
async def get_by_email(session: AsyncSession, email: str) -> User | None:
stmt = select(User).where(User.email == email).limit(1)
return (await session.execute(stmt)).scalar_one_or_none() # 空 → None,多条 → 抛错
async def search_users(session: AsyncSession, keyword: str) -> list[User]:
stmt = select(User).where(User.email.ilike(f"%{keyword}%")).order_by(User.id.desc())
return list((await session.execute(stmt)).scalars().all()) # 注意包一层 list()
async def count_users(session: AsyncSession) -> int:
return (await session.execute(select(func.count()).select_from(User))).scalar_one()
async def create_user(session: AsyncSession, email: str, hashed_password: str) -> User:
user = User(email=email, hashed_password=hashed_password)
session.add(user)
await session.flush() # 发 INSERT,拿到自增 id,但事务未提交
await session.refresh(user) # 取回 server_default 生成的 created_at
return user
async def update_user(session: AsyncSession, user: User, nickname: str) -> User:
user.nickname = nickname # 改属性即可,unit of work 会在 flush 时生成 UPDATE
await session.flush()
return user
async def delete_user(session: AsyncSession, user: User) -> None:
await session.delete(user) # 级联删除已加载的 posts
await session.flush()取值方式必须分清,这是异步 ORM 最容易混淆的一组 API:
| 方法 | 语义 | 结果不匹配时 |
|---|---|---|
result.scalar_one() | 恰好一行一列 | 0 行或多行都抛错 |
result.scalar_one_or_none() | 0 或 1 行 | 多行抛 MultipleResultsFound |
result.scalars().first() | 第一行 | 没数据返回 None,不抛错 |
result.scalars().all() | 全部行的首个实体 | 返回 Sequence,空则 [] |
result.mappings().all() | 每行一个 RowMapping | 适合只取部分列的查询 |
10.8 关系加载:异步下懒加载一定失败
同步 ORM 里 user.posts 会按需发一条 SELECT,这叫懒加载。异步下它必然抛 MissingGreenlet,原因与 10.5 完全相同。
❌ 不推荐:
user = (await session.execute(select(User).where(User.id == 1))).scalar_one()
# 下面这行触发懒加载 → sqlalchemy.exc.MissingGreenlet
for post in user.posts:
print(post.title)✅ 推荐:
from sqlalchemy.orm import joinedload, selectinload
stmt = select(User).where(User.id == 1).options(selectinload(User.posts))
user = (await session.execute(stmt)).scalar_one()
for post in user.posts: # 已在内存,安全访问
print(post.title)
# 多对一用 joinedload,配合 .unique() 去重
stmt = select(Post).options(joinedload(Post.author)).where(Post.id == 1)
post = (await session.execute(stmt)).unique().scalar_one()
print(post.author.email)| 策略 | 生成的 SQL | 适用关系 | 注意 |
|---|---|---|---|
selectinload | 主查询 + 一条 WHERE id IN (...) | 一对多 / 多对多(集合) | 首选,无笛卡尔积 |
joinedload | LEFT OUTER JOIN 一条查询 | 多对一 / 一对一(标量) | 集合关系上用它会让行数重复,必须配 .unique() |
💡 实践建议:把关系显式声明为
relationship(..., lazy="raise"),任何忘记预加载的地方都立刻报错,而不是测试环境「碰巧能用」、上生产才炸。
10.9 Alembic 迁移
Base.metadata.create_all() 只适合脚手架与测试。真实项目表结构会持续演进,必须靠迁移工具记录版本。
uv run alembic init -t async migrations # 异步模板,生成 migrations/ 目录结构为 migrations/env.py(运行时入口,要改)、script.py.mako、versions/(每个迁移一个文件),以及根目录的 alembic.ini。
第一步:alembic.ini 不硬编码连接串(它会进仓库),留空由 env.py 注入:
[alembic]
script_location = migrations
prepend_sys_path = .
sqlalchemy.url =第二步:改 env.py 的三处关键点(节选,asyncio / pool 记得导入):
# migrations/env.py(节选)
from alembic import context
from sqlalchemy.engine import Connection
from sqlalchemy.ext.asyncio import async_engine_from_config
from app.core.config import settings
from app.db.base import Base
# 关键 1:必须 import 所有模型模块,否则 autogenerate 检测不到表
from app.models import post, user # noqa: F401
config = context.config
# 关键 2:连接串从应用配置注入,单一数据源
config.set_main_option("sqlalchemy.url", settings.DATABASE_URL)
# 关键 3:target_metadata 指向 Base.metadata
target_metadata = Base.metadata
def do_run_migrations(connection: Connection) -> None:
context.configure(
connection=connection,
target_metadata=target_metadata,
compare_type=True, # 检测列类型变化
render_as_batch=True, # SQLite 改列的必要开关
)
with context.begin_transaction():
context.run_migrations()
async def run_async_migrations() -> None:
connectable = async_engine_from_config(
config.get_section(config.config_ini_section, {}),
prefix="sqlalchemy.",
poolclass=pool.NullPool,
)
async with connectable.connect() as connection:
await connection.run_sync(do_run_migrations) # 同步迁移逻辑跑进异步连接
await connectable.dispose()第三步:日常命令:
uv run alembic revision --autogenerate -m "create users and posts" # 生成迁移,务必人工 review
uv run alembic upgrade head # 应用全部未执行的迁移
uv run alembic downgrade -1 # 回退一个版本
uv run alembic current # 查看当前版本生成后一定要打开 versions/xxx_*.py 读一遍。autogenerate 只对比模型与数据库的差异:它不知道你要删列、不知道数据如何搬迁,也不检测列重命名(会生成「删掉旧的 + 新建一个空的」)。对生产库做破坏性迁移前先备份。
10.10 开发期建表 vs 生产迁移
# app/main.py
from contextlib import asynccontextmanager
from fastapi import FastAPI
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":
# 异步引擎下 create_all 必须通过 run_sync 调用
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
yield
await engine.dispose() # 优雅关闭:把连接池里的连接全部归还
app = FastAPI(lifespan=lifespan)| 方式 | 建表能力 | 改表能力 | 适用 |
|---|---|---|---|
Base.metadata.create_all | 只创建不存在的表 | 没有(已存在的表一律跳过) | 本地开发、脚手架、CI |
| Alembic 迁移 | 有 | 有,且可回退、可审计 | 生产环境,唯一正确选择 |
⚠️
create_all对已存在的表是静默跳过。改了模型却只重启应用,数据库不会变——这是「本地好好的、换个环境就报no such column」的头号原因。
常见坑与排查
| 现象 | 原因 | 解决 |
|---|---|---|
MissingGreenlet: greenlet_spawn has not been called | 异步上下文里触发了懒加载(关系属性或过期对象) | 查询加 .options(selectinload(...));工厂设 expire_on_commit=False;确认装了 sqlalchemy[asyncio] |
第二个请求开始数据串味、报 InterfaceError | AsyncSession 被当成全局变量跨请求复用 | 会话只在 get_db 里创建、经 Depends 注入,绝不共享实例 |
commit() 之后访问 obj.id 报错 | expire_on_commit=True 使对象过期,访问触发重载 | 设 expire_on_commit=False;需要最新值时 await session.refresh(obj) |
| 数据「改了但没进库」,重启后还原 | 事务没有被提交:只 flush() 没有消费者提交,或压根忘了写操作 | 确认 get_db 在 yield 之后有 await session.commit();Service/Repository 只 flush() |
QueuePool limit of size 5 overflow 10 reached, connection timed out | 会话未释放或泄漏,池被占满 | 确保会话在 async with / yield 依赖内;别在循环里反复 SessionLocal();调整 pool_size / pool_timeout |
SQLite 报 database is locked 或线程相关错误 | 多线程/多事件循环共用一个连接,或写并发过高 | 别跨线程传会话;测试用 StaticPool + check_same_thread=False;生产换 PostgreSQL |
alembic revision --autogenerate 生成空迁移 | env.py 没 import 模型模块,Base.metadata 里没有表 | 在 env.py 显式导入所有模型;或用 app/models/__init__.py 统一导出再导入 |
迁移文件与模型不一致,运行时 no such column | 改了模型却没生成新迁移,或 create_all 静默跳过已存在的表 | 每次模型变更都配一次 revision --autogenerate + upgrade head;禁止用 create_all 改生产表 |
postgresql:// 连接串在 await 时报同步驱动错误 | 未指定异步驱动,SQLAlchemy 选了 psycopg2 | 写成 postgresql+asyncpg://... |
| 池逐渐枯竭、请求卡在等连接 | 会话未走依赖管理,或手工 close() 与 async with 冲突 | 统一用 async with SessionLocal() as session 管理,不要手工关闭 |
本章小结
| 要点 | 说明 |
|---|---|
| 异步驱动 | 开发 sqlite+aiosqlite:///./app.db,生产 postgresql+asyncpg://...,连接串必须带驱动名 |
| 引擎参数 | echo 生产关;pool_pre_ping=True 防死连接;pool_size / max_overflow / pool_recycle 按数据库上限推算 |
| 声明式 | class Base(DeclarativeBase) + Mapped[...] + mapped_column(...),可空性写进类型注解 |
| 关系 | relationship(back_populates=...) 双向写全;ORM 级联与 ondelete="CASCADE" 都配 |
| 会话工厂 | async_sessionmaker(engine, expire_on_commit=False) 是异步下的默认选择 |
| 会话依赖 | async def get_db() -> AsyncIterator[AsyncSession]:yield 前建会话,yield 后提交,异常 rollback(),async with 负责关闭 |
| 事务边界 | 一次请求 = 一个事务,commit() 只在 get_db 里出现一次;Service 与 Repository 只 flush() |
| CRUD | await session.execute(select(...)) → scalars().all() / scalar_one_or_none();add / flush / commit / refresh |
| 关系加载 | 异步禁止隐式懒加载:一对多用 selectinload,多对一用 joinedload + .unique() |
| 迁移 | alembic init -t async migrations,env.py 注入 URL 与 Base.metadata,revision --autogenerate → upgrade head |
| 建表 | create_all 只用于开发(且需 conn.run_sync),生产一律 Alembic |
练习题
写一个
Tag模型,与Post构成多对多关系(中间表post_tags),用selectinload一次性加载某篇Post的全部标签,并说明为什么这里不能用joinedload加载集合。下面这版
get_db刻意去掉了commit()。请写出它的后果,并说明在什么场景下这种「只读会话」反而是更好的选择:
async def get_db() -> AsyncIterator[AsyncSession]:
async with SessionLocal() as session:
yield session用
aiosqlite内存库搭一套测试引擎,要求不同会话能看到同一份数据,写出建表、插入一个用户、再查询该用户的完整片段。给
Post加一个view_count: Mapped[int]列,写出对应的 Alembic 命令,并解释为什么在 SQLite 上必须打开render_as_batch。指出下面代码的三处异步 ORM 问题并修正:
user = (await session.execute(select(User).where(User.id == 7))).scalar_one()
session.commit()
title = user.posts[0].title
return {"email": user.email, "title": title}下一章预告
引擎、会话、模型、迁移都通了,但代码还散在
main.py和几个工具函数里。下一章把它们组织成路由 / Service / Repository 三层结构,让 CRUD 真正工程化。