Python 数据库:SQLAlchemy、事务与迁移

# Python 数据库:SQLAlchemy、事务与迁移

本篇目标

看懂 Python 后端从 Engine、连接池、Session、Repository 到事务和迁移的完整关系,并能判断事务应该放在哪一层。

理解 Engine 与 Session划定事务边界实现 Repository管理 Schema 迁移

# 先记住一句话

Engine 管理数据库方言与连接池,Session 表示一段有状态的工作单元,Repository 封装持久化细节,事务则必须覆盖一个完整业务用例,而不是分散在每条 SQL 后随手提交。

# 四个概念不要混在一起

先按生命周期分成两层:Engine 和 Session 工厂在应用运行期间长期复用,Session 和事务则为每次请求或任务单独创建。

应用启动:创建长期复用的对象
  ├─ AsyncEngine ──管理──→ Connection Pool
  └─ async_sessionmaker ──绑定──→ AsyncEngine

每次请求或任务:创建一次性工作单元
  async_sessionmaker
    └─创建──→ AsyncSession
                ├─按需借用──→ Connection Pool 中的连接
                └─管理──→ Transaction
                            └─结束时──→ 提交或回滚
1
2
3
4
5
6
7
8
9
10

这张图中的箭头表示管理、绑定、创建或借用,不是简单的父子包含关系。SQLAlchemy 的 Session 不是数据库本身,也不是可全局并发共享的连接池。一个 Session 会跟踪加载和修改过的 ORM 对象,并代表一段有顺序的数据库交互;同一个 AsyncSession 不应被多个并发 Task 同时使用。

# 一个可运行的异步数据库层

示例使用 SQLite 和 aiosqlite 让本地容易运行;生产使用 PostgreSQL 时可以把 URL 换成 postgresql+asyncpg://... 并安装 asyncpg。业务代码不应依赖具体驱动。

# 1. 声明依赖

# 文件位置:pyproject.toml
[project]
name = "note-api"
version = "0.1.0"
requires-python = ">=3.12"
dependencies = [
  "sqlalchemy>=2.0,<3.0", # ORM、Engine、Session 与事务 API。
  "aiosqlite>=0.20",     # SQLite 的异步驱动,只用于本地示例。
  "alembic>=1.13",       # 数据库 Schema 迁移工具。
]
1
2
3
4
5
6
7
8
9
10

# 2. 创建 Engine 和 Session 工厂

# 文件位置:src/note_api/database.py
from collections.abc import AsyncIterator

from sqlalchemy.ext.asyncio import (
    AsyncSession,
    async_sessionmaker,
    create_async_engine,
)

DATABASE_URL = "sqlite+aiosqlite:///./notes.db"  # 生产环境应从配置读取。

# Engine 应在应用生命周期内复用,不要每次请求重新创建连接池。
engine = create_async_engine(DATABASE_URL, pool_pre_ping=True)

# expire_on_commit=False 让提交后仍能读取已经加载的属性。
session_factory = async_sessionmaker(engine, expire_on_commit=False)

async def get_session() -> AsyncIterator[AsyncSession]:
    # 每个请求创建独立 Session,退出 async with 时归还连接等资源。
    async with session_factory() as session:
        yield session
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21

pool_pre_ping=True 会在取出连接时检查失效连接,但它不是网络故障的万能重试。连接池大小还要与数据库容量、应用实例数和每个请求的并发查询数量一起计算。

# 3. 定义 ORM 模型

# 文件位置:src/note_api/models.py
from datetime import datetime

from sqlalchemy import DateTime, String, Text, func
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column

# DeclarativeBase 是 SQLAlchemy ORM 模型的共同基类,子类会参与表映射。
class Base(DeclarativeBase):
    """所有 ORM 模型共享的声明式基类。"""

class NoteModel(Base):
    __tablename__ = "notes"  # 显式固定表名,避免类名变动影响数据库。

    # Mapped[str] 表示属性在 Python 中是字符串,并会映射到数据库列。
    # mapped_column 配置实际列;primary_key=True 把 id 设为主键。
    id: Mapped[str] = mapped_column(String(36), primary_key=True)
    # index=True 为常用过滤字段建立索引;String(36) 限制列的最大长度。
    owner_id: Mapped[str] = mapped_column(String(36), index=True)
    title: Mapped[str] = mapped_column(String(200))
    content: Mapped[str] = mapped_column(Text())
    created_at: Mapped[datetime] = mapped_column(
        DateTime(timezone=True),
        server_default=func.now(),  # 时间由数据库生成,多个实例保持一致。
    )
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24

ORM 类型标注同时服务于编辑器和 SQLAlchemy 映射,但外部 HTTP 输入仍需要 Pydantic 校验。数据库模型也不应直接当成公开响应模型,以免将内部字段意外暴露。

# 4. Repository 只负责持久化

# 文件位置:src/note_api/repository.py
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession

from .models import NoteModel

class NoteRepository:
    # Repository(仓储层)是项目代码中的抽象,不是某个固定库;它集中封装数据读写。
    def __init__(self, session: AsyncSession) -> None:
        self._session = session  # Session 生命周期由外层请求或任务控制。

    async def add(self, note: NoteModel) -> None:
        self._session.add(note)
        # flush 把待执行 SQL 发给数据库,但不提交整个事务。
        await self._session.flush()

    async def get_for_owner(self, note_id: str, owner_id: str) -> NoteModel | None:
        # select() 构造查询,where() 添加过滤条件;此时尚未访问数据库。
        statement = select(NoteModel).where(
            NoteModel.id == note_id,
            NoteModel.owner_id == owner_id,  # 权限过滤必须进入数据库条件。
        )
        # scalar() 执行查询并取得第一行的第一个 ORM 对象;查不到时返回 None。
        return await self._session.scalar(statement)

    async def list_for_owner(self, owner_id: str, limit: int) -> list[NoteModel]:
        statement = (
            select(NoteModel)
            .where(NoteModel.owner_id == owner_id)
            .order_by(NoteModel.created_at.desc())
            .limit(limit)
        )
        # scalars() 取得多行中的 ORM 对象序列,最后转成普通列表。
        result = await self._session.scalars(statement)
        return list(result)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35

把 commit() 写进每个 Repository 方法看似方便,却会让 “创建笔记并写审计日志” 无法处于同一个事务。Repository 可以 flush() 获取数据库生成值,是否提交应由更外层的业务用例决定。

# 事务覆盖完整业务用例

# 文件位置:src/note_api/service.py
from dataclasses import dataclass
from uuid import uuid4

from sqlalchemy.ext.asyncio import AsyncSession

from .models import NoteModel
from .repository import NoteRepository

# @dataclass 自动生成初始化等方法;frozen=True 表达命令创建后不应再被修改。
@dataclass(frozen=True)
class CreateNoteCommand:
    owner_id: str
    title: str
    content: str

class NoteService:
    def __init__(self, session: AsyncSession) -> None:
        self._session = session
        self._notes = NoteRepository(session)

    async def create(self, command: CreateNoteCommand) -> NoteModel:
        note = NoteModel(
            # uuid4() 生成随机 UUID;str() 把它转成数据库字段需要的字符串。
            id=str(uuid4()),
            owner_id=command.owner_id,
            title=command.title,
            content=command.content,
        )

        # begin 在成功退出时提交;出现异常时自动回滚整个用例。
        async with self._session.begin():
            await self._notes.add(note)
            # 如果还要写审计表,应在同一事务作用域内调用对应 Repository。

        return note

    async def get_for_owner(self, note_id: str, owner_id: str) -> NoteModel | None:
        # 查询方法继续复用 Repository 中的权限条件,入口层不自行拼 SQL。
        return await self._notes.get_for_owner(note_id, owner_id)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40

不要在事务中执行耗时模型调用或等待用户操作,否则数据库连接和锁会被长时间占用。常见顺序是先完成外部计算,再开启短事务写结果;若必须保证 “数据库状态与消息发布” 一致,需要 Outbox 等专门模式,而不是让数据库事务跨网络服务。

# 连接池不是越大越好

假设数据库最多允许 100 条业务连接,应用有 5 个实例,如果每个实例都配置 30 个连接,理论峰值会达到 150,数据库可能在扩容时直接被打满。

数据库可用连接预算
└─ 除以应用最大实例数
   └─ 再预留管理、迁移和故障切换空间
      └─ 得到单实例连接池上限
1
2
3
4

连接池过小会排队,过大则把压力直接推给数据库。应观测连接等待时间、活跃连接、慢查询和事务时长,再调整池参数。

# Schema 变化必须用迁移记录

Alembic (opens new window) 是 SQLAlchemy 常用的数据库 Schema 迁移工具。它把建表、加列和修改索引等结构变化保存成带版本的迁移脚本,可以简单理解为数据库表结构的版本管理工具。

开发环境用 Base.metadata.create_all() 可以快速建空表,但它不会安全地演进已有生产表。正式项目使用 Alembic 生成并提交版本化迁移。

# 在项目根目录初始化迁移目录;整个项目通常只执行一次。
alembic init migrations

# 比较 ORM metadata 和当前数据库,生成候选迁移文件。
alembic revision --autogenerate -m "create notes table"

# 人工审查迁移脚本后,将数据库升级到最新版本。
alembic upgrade head
1
2
3
4
5
6
7
8
# 文件位置:migrations/env.py
# 这是 Alembic 生成文件中需要修改的关键片段,不是完整 env.py。
from note_api.models import Base

# Alembic 使用这份 metadata 与数据库现状比较。
target_metadata = Base.metadata
1
2
3
4
5
6

自动生成只提供候选迁移,仍需人工检查数据回填、列重命名、锁表风险和回滚策略。部署时也不要让每个 Web Worker 同时自动跑迁移,应由发布流程中的单独步骤负责。

# 高频面试题与回答

1. Engine、连接池和 Session 分别是什么?参考答案

Engine 是应用访问某类数据库的入口,通常持有连接池;连接池复用有限的真实连接;Session 是一次工作单元,跟踪 ORM 对象并管理事务状态。Engine 通常全应用复用,Session 通常按请求或任务创建,不能让多个并发 Task 共用一个 AsyncSession。

2. 为什么 Repository 不应该随手 commit?参考答案

事务应覆盖完整业务用例。如果每个 Repository 方法自行提交,多个持久化动作就无法原子成功或回滚。Repository 负责查询和写入,Service 或 Unit of Work 根据业务边界决定提交。

3. Alembic 自动生成的迁移能直接上线吗?参考答案

不能默认直接上线。自动生成只能比较部分 Schema 差异,列重命名可能被识别成删除再新增,也不了解数据回填和大表锁风险。迁移文件应纳入版本控制,经过人工审查、测试和发布策略验证。

4. 数据库连接池是不是越大越好?参考答案

不是。连接池过小会让请求等待,过大则会让所有应用实例一起压垮数据库,并增加内存和上下文切换。应根据实例数、数据库连接上限、查询耗时和峰值并发设定池大小与溢出上限,再通过等待时间、超时和数据库负载调整。

# 接下来学什么

下一篇学习 FastAPI 生产接口,把 Session、身份和应用生命周期通过依赖注入接入 HTTP 请求。

# 参考资料

上次更新时间: 2026年09月10日 00:03:52