Skip to content

第 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 卡住,异步的并发能力就白给了。

组合驱动示例连接串适用场景
同步 ORMpsycopg2 / sqlite3postgresql://...传统 WSGI 应用、Celery 任务
同步 ORM + def 端点同上同上过渡期,FastAPI 会把 def 丢进线程池
异步 ORMasyncpg / aiosqlitepostgresql+asyncpg://...本教程统一采用

选它的理由:异步会话能直接用在 async def 端点、WebSocket、后台任务、httpx 调用之后,不必来回切线程池;AsyncSession / async_sessionmaker / create_async_engine 都是官方原生 API;连接池行为完全显式,没有隐式线程抢占。

环境数据库驱动包连接串
开发SQLite 文件aiosqlitesqlite+aiosqlite:///./app.db
测试SQLite 内存aiosqlitesqlite+aiosqlite:///:memory:
生产PostgreSQLasyncpgpostgresql+asyncpg://user:pwd@host:5432/db
bash
# 基础三件套:异步扩展 + 开发用 SQLite 驱动 + 迁移工具
uv add "sqlalchemy[asyncio]" aiosqlite alembic

# 生产环境再加 PostgreSQL 异步驱动
uv add asyncpg

sqlalchemy[asyncio] 的 extra 会带上 greenlet。它是异步 ORM 的硬依赖,缺它不会在导入时报错,而是在第一次查库时抛 MissingGreenlet——典型「装了但装漏了」的坑。

连接串格式为 方言+驱动://用户:密码@主机:端口/库名:

text
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 创建异步引擎 ​

python
# 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:

python
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(...):

python
# 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 旧写法):

python
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 写法):

python
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 ​

python
# 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,
    )
python
# 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")

三处细节:

  1. back_populates 必须双向写全。只写一边,另一边不会同步,user.posts 与 post.author 会指向不同的 Python 对象。
  2. cascade="all, delete-orphan" 是 ORM 层行为:删除 user 时,会话自动对已加载的 post 发 DELETE;把 post 从 user.posts 移除,它也会被删除。
  3. 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 ​

python
# 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 手上只有异步驱动,于是直接抛错:

text
sqlalchemy.exc.MissingGreenlet: greenlet_spawn has not been called;
can't call await_only() here. Was IO attempted in an unexpected place?

❌ 不推荐:

python
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         # 路由层序列化时同样炸

✅ 推荐:

python
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 之后清理。数据库会话是它最经典的用例。

python
# 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()   # 任何异常 → 整体回滚
            raise
  • async with SessionLocal() 已负责关闭,不要手写 session.close()。
  • commit() 只在这里出现一次:整个请求是一个事务,yield 之后的代码在路由返回后执行,成功则提交、出错则回滚。这样就不会出现「前半段写进去了、后半段失败」的半提交。这一决策的完整取舍见第 11 章 11.6。
  • rollback() 是必需的兜底,否则异常路径会把一个未结束的事务留在连接上,连接归还连接池时状态是脏的。
  • 绝不要把 session 存成模块级全局变量复用。AsyncSession 既不线程安全也不协程安全,跨请求共享会导致 InterfaceError 与数据串味。

10.7 异步 CRUD 标准动作 ​

python
# 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 完全相同。

❌ 不推荐:

python
user = (await session.execute(select(User).where(User.id == 1))).scalar_one()

# 下面这行触发懒加载 → sqlalchemy.exc.MissingGreenlet
for post in user.posts:
    print(post.title)

✅ 推荐:

python
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 (...)一对多 / 多对多(集合)首选,无笛卡尔积
joinedloadLEFT OUTER JOIN 一条查询多对一 / 一对一(标量)集合关系上用它会让行数重复,必须配 .unique()

💡 实践建议:把关系显式声明为 relationship(..., lazy="raise"),任何忘记预加载的地方都立刻报错,而不是测试环境「碰巧能用」、上生产才炸。


10.9 Alembic 迁移 ​

Base.metadata.create_all() 只适合脚手架与测试。真实项目表结构会持续演进,必须靠迁移工具记录版本。

bash
uv run alembic init -t async migrations   # 异步模板,生成 migrations/ 目录

结构为 migrations/env.py(运行时入口,要改)、script.py.mako、versions/(每个迁移一个文件),以及根目录的 alembic.ini。

第一步:alembic.ini 不硬编码连接串(它会进仓库),留空由 env.py 注入:

ini
[alembic]
script_location = migrations
prepend_sys_path = .
sqlalchemy.url =

第二步:改 env.py 的三处关键点(节选,asyncio / pool 记得导入):

python
# 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()

第三步:日常命令:

bash
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 生产迁移 ​

python
# 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]
第二个请求开始数据串味、报 InterfaceErrorAsyncSession 被当成全局变量跨请求复用会话只在 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()
CRUDawait 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

练习题 ​

  1. 写一个 Tag 模型,与 Post 构成多对多关系(中间表 post_tags),用 selectinload 一次性加载某篇 Post 的全部标签,并说明为什么这里不能用 joinedload 加载集合。

  2. 下面这版 get_db 刻意去掉了 commit()。请写出它的后果,并说明在什么场景下这种「只读会话」反而是更好的选择:

python
async def get_db() -> AsyncIterator[AsyncSession]:
    async with SessionLocal() as session:
        yield session
  1. 用 aiosqlite 内存库搭一套测试引擎,要求不同会话能看到同一份数据,写出建表、插入一个用户、再查询该用户的完整片段。

  2. 给 Post 加一个 view_count: Mapped[int] 列,写出对应的 Alembic 命令,并解释为什么在 SQLite 上必须打开 render_as_batch。

  3. 指出下面代码的三处异步 ORM 问题并修正:

python
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 真正工程化。

👉 第 11 章:分层架构与 CRUD 工程化

📖本文阅读--次|📊全站访问--次|👥访客--人