SQLAlchemy 完全指南

SQLAlchemy 由两层组成: 通常使用 ORM 层即可,Core 层在需要批量操作或高性能场景下使用。 SQLAlchemy 2.0 是一次重大版本升级(2023年发布),主要变化: 本文重点介绍 2.x 写法。 常用连接字符串格式: 参数与 create_engine() 相同,但 URL 需使用异步驱动(asyncpg、aiomysql、aiosqlite)。 外键定义: 注意:scoped_session 基于线程本地存储,不适用于异步环境。异步请使用 async_sessionmaker + 依赖注入。 selectinload/joine

分享

官方文档:https://docs.sqlalchemy.org/
最后更新:2026-03-05


1. 基础概念

ORM vs Core

SQLAlchemy 由两层组成:

层次 名称 说明
高层 ORM(对象关系映射) 将 Python 类映射到数据库表,通过对象操作数据
低层 Core 基于 SQL 表达式语言,更接近原生 SQL,性能更高

通常使用 ORM 层即可,Core 层在需要批量操作或高性能场景下使用。

SQLAlchemy 2.x 的主要变化

SQLAlchemy 2.0 是一次重大版本升级(2023年发布),主要变化:

特性 1.x 旧写法 2.x 新写法
模型基类 Base = declarative_base() class Base(DeclarativeBase): pass
字段定义 Column(Integer, ...) mapped_column(int, ...) 配合 Mapped[int]
查询 session.query(User).filter(...) select(User).where(...)
执行查询 session.query(User).all() session.scalars(select(User)).all()
异步支持 需要 sqlalchemy-aio 等第三方 原生 create_async_engine + AsyncSession
类型注解 不强制 完整的 Mapped[T] 类型提示

本文重点介绍 2.x 写法。


2. 连接与引擎

create_engine() 参数表

from sqlalchemy import create_engine

# 同步引擎
engine = create_engine(
    "postgresql+psycopg2://user:password@localhost:5432/mydb",
    pool_size=5,
    max_overflow=10,
    pool_timeout=30,
    pool_recycle=3600,
    echo=False,
)
参数 类型 默认值 说明
url str / URL 必填 数据库连接字符串
echo bool / str False True 打印所有 SQL,"debug" 打印参数
pool_size int 5 连接池保持的连接数
max_overflow int 10 超出 pool_size 时最多额外创建的连接数
pool_timeout float 30 等待连接超时秒数
pool_recycle float -1 连接重用最大秒数,防止数据库断开,建议 3600
pool_pre_ping bool False 每次使用前检测连接是否存活,防止断连
connect_args dict {} 传递给底层驱动的额外参数
execution_options dict {} 执行选项
isolation_level str None 隔离级别,如 "READ COMMITTED"

常用连接字符串格式:

# PostgreSQL(同步)
postgresql+psycopg2://user:pass@host:5432/dbname

# PostgreSQL(异步)
postgresql+asyncpg://user:pass@host:5432/dbname

# MySQL(同步)
mysql+pymysql://user:pass@host:3306/dbname?charset=utf8mb4

# MySQL(异步)
mysql+aiomysql://user:pass@host:3306/dbname

# SQLite(同步)
sqlite:///./local.db
sqlite:///:memory:

# SQLite(异步)
sqlite+aiosqlite:///./local.db

create_async_engine()

from sqlalchemy.ext.asyncio import create_async_engine

engine = create_async_engine(
    "postgresql+asyncpg://user:password@localhost/mydb",
    pool_size=10,
    max_overflow=20,
    pool_pre_ping=True,    # 推荐开启,防止连接断开
    echo=False,
)

参数与 create_engine() 相同,但 URL 需使用异步驱动(asyncpg、aiomysql、aiosqlite)。

连接池配置建议

# 生产环境推荐配置
engine = create_async_engine(
    DATABASE_URL,
    pool_size=10,           # 核心连接数,根据服务器核心数调整
    max_overflow=20,        # 突发流量时额外连接
    pool_timeout=30,        # 等待超时
    pool_recycle=1800,      # 30分钟回收,防止 MySQL 的 wait_timeout 断开
    pool_pre_ping=True,     # 使用前检测,推荐开启
)

3. 模型定义

DeclarativeBase

from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship
from sqlalchemy import String, Integer, DateTime, func
from datetime import datetime

# 2.x 新写法:继承 DeclarativeBase
class Base(DeclarativeBase):
    pass

class User(Base):
    __tablename__ = "users"

    id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
    username: Mapped[str] = mapped_column(String(50), unique=True, nullable=False)
    email: Mapped[str] = mapped_column(String(100), unique=True)
    password_hash: Mapped[str] = mapped_column(String(255))
    is_active: Mapped[bool] = mapped_column(default=True)
    created_at: Mapped[datetime] = mapped_column(
        DateTime(timezone=True),
        server_default=func.now(),   # 数据库层面的默认值
    )
    updated_at: Mapped[datetime | None] = mapped_column(
        DateTime(timezone=True),
        onupdate=func.now(),         # 更新时自动设置
    )

    # 关系
    posts: Mapped[list["Post"]] = relationship(back_populates="author")

mapped_column() 参数表

参数 类型 默认值 说明
type_ TypeEngine 从注解推断 字段类型,如 String(50)
primary_key bool False 是否为主键
nullable bool True 是否允许 NULL(Mapped[T] 非 Optional 时自动推断为 False)
unique bool False 是否唯一约束
index bool False 是否创建索引
default Any None Python 层面的默认值(insert 时)
server_default str / FetchedValue None 数据库层面的默认值(SQL 表达式)
onupdate Any None 更新时自动设置的值
autoincrement bool / str "auto" 是否自增
comment str None 字段注释
name str None 数据库列名(默认与属性名相同)
foreign_keys set None 外键约束
init bool True 是否包含在 __init__ 中(dataclass 集成)

常用数据类型

from sqlalchemy import (
    Integer, BigInteger, SmallInteger,
    Float, Numeric,
    String, Text, Unicode,
    Boolean,
    Date, DateTime, Time, Interval,
    JSON, ARRAY,
    Enum,
    LargeBinary,
)

class Example(Base):
    __tablename__ = "examples"

    id: Mapped[int] = mapped_column(Integer, primary_key=True)
    big_id: Mapped[int] = mapped_column(BigInteger)
    price: Mapped[float] = mapped_column(Numeric(10, 2))  # 精确小数,10位总长,2位小数
    name: Mapped[str] = mapped_column(String(100))
    content: Mapped[str] = mapped_column(Text)
    flag: Mapped[bool] = mapped_column(Boolean)
    created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True))
    data: Mapped[dict] = mapped_column(JSON)
    status: Mapped[str] = mapped_column(Enum("active", "inactive", name="status_enum"))

relationship() 参数表

from sqlalchemy.orm import relationship, Mapped

class Post(Base):
    __tablename__ = "posts"

    id: Mapped[int] = mapped_column(primary_key=True)
    author_id: Mapped[int] = mapped_column(ForeignKey("users.id"))
    title: Mapped[str] = mapped_column(String(200))

    author: Mapped["User"] = relationship(
        back_populates="posts",
        lazy="select",
    )
参数 类型 默认值 说明
back_populates str None 反向关系属性名,双向关系时使用
backref str None 自动创建反向关系(简单场景)
lazy str "select" 加载策略,见第 7 节
foreign_keys list None 有歧义时指定外键列
primaryjoin str None 自定义 join 条件
cascade str "save-update, merge" 级联操作,如 "all, delete-orphan"
uselist bool 推断 False 表示一对一关系
order_by Column None 默认排序
secondary Table None 多对多中间表
overlaps str None 抑制重叠警告

外键定义:

from sqlalchemy import ForeignKey

class Post(Base):
    __tablename__ = "posts"
    author_id: Mapped[int] = mapped_column(ForeignKey("users.id", ondelete="CASCADE"))

4. Session 使用

Session / AsyncSession

from sqlalchemy.orm import Session
from sqlalchemy.ext.asyncio import AsyncSession

# 同步 Session
with Session(engine) as session:
    user = session.get(User, 1)
    print(user.username)

# 异步 AsyncSession
async with AsyncSession(engine) as session:
    user = await session.get(User, 1)
    print(user.username)

sessionmaker / async_sessionmaker

from sqlalchemy.orm import sessionmaker
from sqlalchemy.ext.asyncio import async_sessionmaker

# 同步
SyncSessionLocal = sessionmaker(
    bind=engine,
    autocommit=False,    # 不自动提交
    autoflush=True,      # 查询前自动 flush
    expire_on_commit=True,  # commit 后实例过期(下次访问重新查询)
)

# 异步(推荐)
AsyncSessionLocal = async_sessionmaker(
    engine,
    expire_on_commit=False,  # 异步场景推荐 False,避免 lazy load 问题
    autoflush=True,
)
参数 类型 默认值 说明
bind Engine None 绑定的引擎
autocommit bool False 是否自动提交(不推荐开启)
autoflush bool True 查询前是否自动 flush pending 的更改
expire_on_commit bool True commit 后实例属性是否过期
class_ type Session 使用的 Session 类

上下文管理器

# 推荐写法:自动处理提交/回滚/关闭
async def create_user(username: str):
    async with AsyncSessionLocal() as session:
        async with session.begin():            # 开始事务,出错自动回滚
            user = User(username=username)
            session.add(user)
        # session.begin() 结束时自动 commit
    # async with 结束时自动 close

scoped_session(线程安全的 Session 注册表)

from sqlalchemy.orm import scoped_session, sessionmaker

session_factory = sessionmaker(bind=engine)
ScopedSession = scoped_session(session_factory)

# 同一线程内多次调用返回同一 session
session = ScopedSession()
session.add(user)
session.commit()
ScopedSession.remove()  # 使用完毕后释放

注意scoped_session 基于线程本地存储,不适用于异步环境。异步请使用 async_sessionmaker + 依赖注入。


5. 查询(2.x select 风格)

基础查询

from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession

async def get_all_users(session: AsyncSession):
    # 查询所有用户
    stmt = select(User)
    result = await session.execute(stmt)
    users = result.scalars().all()         # 返回 User 实例列表
    return users

async def get_user_by_id(session: AsyncSession, user_id: int):
    # 主键查询(推荐)
    user = await session.get(User, user_id)
    return user

async def get_one_user(session: AsyncSession, username: str):
    stmt = select(User).where(User.username == username)
    user = await session.scalar(stmt)      # 返回第一个结果或 None
    return user

where() 条件

from sqlalchemy import select, and_, or_, not_, in_

stmt = (
    select(User)
    .where(User.is_active == True)                   # 等于
    .where(User.age >= 18)                           # 比较
    .where(User.username.like("%admin%"))             # LIKE
    .where(User.username.ilike("%Admin%"))            # 大小写不敏感 LIKE
    .where(User.email.in_(["[email protected]", "[email protected]"]))   # IN
    .where(User.deleted_at.is_(None))                # IS NULL
    .where(User.deleted_at.is_not(None))             # IS NOT NULL
)

# AND / OR / NOT
stmt = select(User).where(
    and_(
        User.is_active == True,
        or_(
            User.role == "admin",
            User.role == "superuser",
        )
    )
)

# between
stmt = select(User).where(User.age.between(18, 30))

join() / outerjoin()

# INNER JOIN
stmt = (
    select(Post, User)
    .join(User, Post.author_id == User.id)
    .where(User.is_active == True)
)

# LEFT OUTER JOIN
stmt = (
    select(User, Post)
    .outerjoin(Post, User.id == Post.author_id)
)

# 通过关系自动 join(有 relationship 定义时)
stmt = select(Post).join(Post.author).where(User.is_active == True)

order_by() / limit() / offset()

from sqlalchemy import desc, asc

stmt = (
    select(User)
    .where(User.is_active == True)
    .order_by(desc(User.created_at), asc(User.username))  # 多字段排序
    .limit(20)
    .offset(40)                                            # 第 3 页(每页 20 条)
)

获取结果的方法

result = await session.execute(stmt)

result.scalars().all()      # 返回 list[Model],多列时取第一列
result.all()                # 返回 list[Row],每行是元组
result.scalar_one()         # 返回恰好一条,0 条或多条都抛异常
result.scalar_one_or_none() # 返回一条或 None,多条抛异常
result.first()              # 返回第一条 Row 或 None
result.one()                # 返回恰好一条 Row,否则抛异常
result.one_or_none()        # 返回一条 Row 或 None,多条抛异常

# 快捷写法
await session.scalar(stmt)  # 等同于 execute(stmt).scalar_one_or_none()
await session.scalars(stmt) # 等同于 execute(stmt).scalars()

聚合与分组

from sqlalchemy import func, select

# COUNT
stmt = select(func.count(User.id)).where(User.is_active == True)
count = await session.scalar(stmt)

# GROUP BY + HAVING
stmt = (
    select(User.role, func.count(User.id).label("user_count"))
    .group_by(User.role)
    .having(func.count(User.id) > 5)
    .order_by(desc("user_count"))
)

# 其他聚合函数
func.sum(Order.amount)
func.avg(Order.amount)
func.max(Order.amount)
func.min(Order.amount)

6. 增删改

新增

# 新增单条
async def create_user(session: AsyncSession, username: str, email: str):
    user = User(username=username, email=email)
    session.add(user)
    await session.commit()
    await session.refresh(user)  # 刷新,获取数据库生成的字段(如 id、created_at)
    return user

# 新增多条
async def create_users_bulk(session: AsyncSession, users_data: list[dict]):
    users = [User(**data) for data in users_data]
    session.add_all(users)
    await session.commit()

删除

from sqlalchemy import delete

# 通过对象删除
async def delete_user(session: AsyncSession, user_id: int):
    user = await session.get(User, user_id)
    if user:
        await session.delete(user)
        await session.commit()

# 批量删除(不加载到内存,性能更好)
async def delete_inactive_users(session: AsyncSession):
    stmt = delete(User).where(User.is_active == False)
    await session.execute(stmt)
    await session.commit()

更新

from sqlalchemy import update

# 通过对象更新
async def update_user_email(session: AsyncSession, user_id: int, new_email: str):
    user = await session.get(User, user_id)
    if user:
        user.email = new_email
        await session.commit()
        await session.refresh(user)
    return user

# 批量更新(不加载到内存)
async def activate_all_users(session: AsyncSession):
    stmt = (
        update(User)
        .where(User.is_active == False)
        .values(is_active=True)
    )
    await session.execute(stmt)
    await session.commit()

bulk_insert_mappings()(高性能批量插入)

# 2.x 推荐方式:直接用 insert() 语句
from sqlalchemy.dialects.postgresql import insert

async def bulk_insert_users(session: AsyncSession, users_data: list[dict]):
    stmt = insert(User).values(users_data)
    await session.execute(stmt)
    await session.commit()

# 或者使用 Core 层的批量插入
await session.execute(
    User.__table__.insert(),
    users_data,  # list of dicts
)
await session.commit()

7. 关系查询

加载策略概览

策略 触发时机 生成 SQL 适用场景
lazy="select"(默认) 访问属性时 额外一条 SELECT 简单同步场景
lazy="joined" 主查询时 JOIN 一对一或小型一对多
lazy="subquery" 主查询时 子查询 SELECT 一对多
lazy="selectin" 主查询后 IN 子查询 一对多,异步推荐
lazy="dynamic" 返回 Query 对象 手动执行 已废弃,2.x 不推荐
lazy="raise" 访问时直接报错 防止意外 lazy load
lazy="noload" 不加载 不需要关系数据时

显式预加载(查询时指定)

from sqlalchemy.orm import selectinload, joinedload, subqueryload

# selectinload:推荐,适合一对多,异步友好
stmt = (
    select(User)
    .options(selectinload(User.posts))  # 预加载 posts
    .where(User.is_active == True)
)
users = await session.scalars(stmt)

# joinedload:适合一对一或多对一
stmt = (
    select(Post)
    .options(joinedload(Post.author))   # 用 JOIN 加载 author
)

# 嵌套预加载
stmt = (
    select(User)
    .options(
        selectinload(User.posts).selectinload(Post.comments)  # 链式预加载
    )
)

# 多个关系同时预加载
stmt = (
    select(User)
    .options(
        selectinload(User.posts),
        selectinload(User.orders),
    )
)

selectinload/joinedload/subqueryload 参数表:

参数 类型 说明
attr 关系属性 要预加载的关系属性,如 User.posts
链式调用 预加载选项 可继续调用 .selectinload() 等实现嵌套

8. 事务

with session.begin()

# 方式 1:使用 begin() 上下文管理器
async def transfer_money(session: AsyncSession, from_id: int, to_id: int, amount: float):
    async with session.begin():
        from_account = await session.get(Account, from_id)
        to_account = await session.get(Account, to_id)
        from_account.balance -= amount
        to_account.balance += amount
        # 出错时自动 rollback,成功时自动 commit

# 方式 2:手动控制
async def manual_transaction(session: AsyncSession):
    try:
        user = User(username="test")
        session.add(user)
        await session.flush()      # 发送 SQL 但不提交(可获取 id)
        print(user.id)             # 此时 id 已可用
        await session.commit()
    except Exception:
        await session.rollback()
        raise

嵌套事务(SAVEPOINT)

async def nested_transaction(session: AsyncSession):
    async with session.begin():
        user = User(username="outer")
        session.add(user)

        # 嵌套事务使用 SAVEPOINT
        async with session.begin_nested():
            try:
                risky_user = User(username="risky")
                session.add(risky_user)
                await session.flush()
                # 如果这里抛异常,只回滚到 SAVEPOINT,外层事务继续
            except Exception:
                # SAVEPOINT 会自动回滚
                pass

        # 外层事务继续,user 仍然会被提交

9. Alembic 数据库迁移

安装与初始化

pip install alembic

# 初始化(在项目根目录执行)
alembic init alembic

初始化后生成:

alembic/
├── env.py           # 迁移环境配置
├── versions/        # 迁移脚本存放目录
└── script.py.mako   # 迁移脚本模板
alembic.ini          # Alembic 配置文件

配置 env.py

# alembic/env.py 关键配置
from myapp.models import Base  # 导入模型 Base
from myapp.config import settings

# 设置数据库 URL
config.set_main_option("sqlalchemy.url", settings.database_url)

# 设置 target_metadata(用于自动生成迁移)
target_metadata = Base.metadata

常用命令表

命令 说明
alembic init alembic 初始化 Alembic 目录
alembic revision -m "描述" 创建空迁移脚本
alembic revision --autogenerate -m "描述" 自动检测模型变化生成迁移脚本
alembic upgrade head 执行所有未应用的迁移(升级到最新)
alembic upgrade +1 执行下一个迁移
alembic downgrade -1 回滚最近一个迁移
alembic downgrade base 回滚所有迁移
alembic current 查看当前数据库版本
alembic history 查看迁移历史
alembic show <revision> 查看指定迁移内容

迁移脚本示例

# alembic/versions/xxxx_add_users_table.py
"""add users table

Revision ID: xxxx
Revises:
Create Date: 2026-03-05
"""
from alembic import op
import sqlalchemy as sa

revision = "xxxx"
down_revision = None   # 上一个迁移的 ID

def upgrade() -> None:
    op.create_table(
        "users",
        sa.Column("id", sa.Integer(), primary_key=True),
        sa.Column("username", sa.String(50), nullable=False),
        sa.Column("email", sa.String(100), unique=True),
    )
    op.create_index("ix_users_username", "users", ["username"], unique=True)

def downgrade() -> None:
    op.drop_index("ix_users_username", "users")
    op.drop_table("users")

异步 Alembic 配置

# alembic/env.py(异步版本)
import asyncio
from sqlalchemy.ext.asyncio import create_async_engine

def run_migrations_online():
    connectable = create_async_engine(config.get_main_option("sqlalchemy.url"))

    async def run_async_migrations():
        async with connectable.connect() as connection:
            await connection.run_sync(do_run_migrations)

    asyncio.run(run_async_migrations())

10. 与 FastAPI 集成

async session 依赖

# database.py
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession
from typing import AsyncGenerator

DATABASE_URL = "postgresql+asyncpg://user:pass@localhost/mydb"

engine = create_async_engine(DATABASE_URL, pool_pre_ping=True)
AsyncSessionLocal = async_sessionmaker(engine, expire_on_commit=False)

async def get_db() -> AsyncGenerator[AsyncSession, None]:
    async with AsyncSessionLocal() as session:
        yield session

# 在路由中使用
from fastapi import Depends

@app.get("/users/{user_id}")
async def get_user(user_id: int, db: AsyncSession = Depends(get_db)):
    user = await db.get(User, user_id)
    if not user:
        raise HTTPException(status_code=404, detail="用户不存在")
    return user

lifespan 创建表

# main.py
from contextlib import asynccontextmanager
from fastapi import FastAPI
from database import engine, Base

@asynccontextmanager
async def lifespan(app: FastAPI):
    # 应用启动:创建数据库表
    async with engine.begin() as conn:
        await conn.run_sync(Base.metadata.create_all)
    yield
    # 应用关闭:释放连接池
    await engine.dispose()

app = FastAPI(lifespan=lifespan)

详见 FastAPI完全指南


11. 高级特性

混入(Mixin)

from sqlalchemy.orm import Mapped, mapped_column
from sqlalchemy import DateTime, func
from datetime import datetime

class TimestampMixin:
    """为模型自动添加 created_at 和 updated_at 字段"""
    created_at: Mapped[datetime] = mapped_column(
        DateTime(timezone=True),
        server_default=func.now(),
    )
    updated_at: Mapped[datetime | None] = mapped_column(
        DateTime(timezone=True),
        onupdate=func.now(),
    )

class SoftDeleteMixin:
    """软删除支持"""
    deleted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))

    @property
    def is_deleted(self) -> bool:
        return self.deleted_at is not None

# 使用混入
class User(Base, TimestampMixin, SoftDeleteMixin):
    __tablename__ = "users"
    id: Mapped[int] = mapped_column(primary_key=True)
    username: Mapped[str] = mapped_column(String(50))

事件监听

from sqlalchemy import event
from sqlalchemy.orm import Session

# 监听 ORM 事件
@event.listens_for(User, "before_insert")
def user_before_insert(mapper, connection, target):
    # target 是即将插入的 User 实例
    if target.username:
        target.username = target.username.lower()  # 自动转小写

# 监听 Session 事件
@event.listens_for(Session, "after_commit")
def after_commit(session):
    print("事务已提交")

混合属性(hybrid_property)

from sqlalchemy.ext.hybrid import hybrid_property
from sqlalchemy import select, func

class User(Base):
    __tablename__ = "users"

    id: Mapped[int] = mapped_column(primary_key=True)
    first_name: Mapped[str] = mapped_column(String(50))
    last_name: Mapped[str] = mapped_column(String(50))

    @hybrid_property
    def full_name(self) -> str:
        """Python 层:通过实例访问"""
        return f"{self.first_name} {self.last_name}"

    @full_name.expression
    @classmethod
    def full_name(cls):
        """SQL 层:用于 where() 查询"""
        return func.concat(cls.first_name, " ", cls.last_name)

# 使用
user = await session.get(User, 1)
print(user.full_name)  # Python 层

stmt = select(User).where(User.full_name == "张 三")  # SQL 层

12. 最佳实践

Repository 模式

# repositories/user_repository.py
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession

class UserRepository:
    def __init__(self, session: AsyncSession):
        self.session = session

    async def get_by_id(self, user_id: int) -> User | None:
        return await self.session.get(User, user_id)

    async def get_by_username(self, username: str) -> User | None:
        stmt = select(User).where(User.username == username)
        return await self.session.scalar(stmt)

    async def create(self, username: str, email: str) -> User:
        user = User(username=username, email=email)
        self.session.add(user)
        await self.session.flush()  # 获取 id 但不提交
        return user

    async def get_active_users(self, skip: int = 0, limit: int = 10) -> list[User]:
        stmt = (
            select(User)
            .where(User.is_active == True)
            .offset(skip)
            .limit(limit)
        )
        return list(await self.session.scalars(stmt))

# 在 FastAPI 中使用
def get_user_repository(db: AsyncSession = Depends(get_db)) -> UserRepository:
    return UserRepository(db)

@app.get("/users/{user_id}")
async def get_user(
    user_id: int,
    repo: UserRepository = Depends(get_user_repository),
):
    user = await repo.get_by_id(user_id)
    if not user:
        raise HTTPException(status_code=404, detail="用户不存在")
    return user

async 最佳配置

# database.py 生产环境推荐配置
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker

engine = create_async_engine(
    DATABASE_URL,
    pool_size=10,             # 根据并发量调整
    max_overflow=20,
    pool_timeout=30,
    pool_recycle=1800,        # 防止 MySQL wait_timeout 断开
    pool_pre_ping=True,       # 使用连接前检测有效性
    echo=False,               # 生产环境关闭
)

# expire_on_commit=False 非常重要
# 若为 True,commit 后属性过期,在异步场景访问会触发懒加载(报错)
AsyncSessionLocal = async_sessionmaker(
    engine,
    expire_on_commit=False,
    autoflush=True,
)

13. 常见陷阱

lazy load 在异步中的问题

# 陷阱:异步 session 中访问懒加载属性会报错
async def bad_example(session: AsyncSession):
    user = await session.get(User, 1)
    print(user.posts)  # 错误!触发懒加载,但在异步环境中无法隐式执行 IO

# 正确做法 1:使用 selectinload 预加载
async def good_example_1(session: AsyncSession):
    stmt = select(User).options(selectinload(User.posts)).where(User.id == 1)
    user = await session.scalar(stmt)
    print(user.posts)  # 已预加载,安全访问

# 正确做法 2:在 relationship 上设置 lazy="raise"(提前发现问题)
class User(Base):
    posts: Mapped[list["Post"]] = relationship(lazy="raise")
    # 任何未预加载的访问都会抛出 MissingGreenlet 异常

# 正确做法 3:在 session 中显式加载
from sqlalchemy.orm import selectinload
await session.refresh(user, ["posts"])  # 显式加载指定关系

DetachedInstanceError

# 陷阱:在 session 关闭后访问属性
async def bad_example():
    async with AsyncSessionLocal() as session:
        user = await session.get(User, 1)
    # session 已关闭
    print(user.username)  # 如果 expire_on_commit=True,会报 DetachedInstanceError

# 原因:session 关闭或 commit 后,若 expire_on_commit=True,
# 实例属性会被标记为过期,再次访问会触发 lazy load,
# 但此时 session 已关闭无法加载。

# 解决方案 1:设置 expire_on_commit=False(异步推荐)
AsyncSessionLocal = async_sessionmaker(engine, expire_on_commit=False)

# 解决方案 2:在 session 内完成所有访问(包括序列化)
async def good_example():
    async with AsyncSessionLocal() as session:
        user = await session.get(User, 1)
        result = {"username": user.username, "email": user.email}  # 在 session 内访问
    return result  # session 外只使用普通字典

# 解决方案 3:使用 response_model(FastAPI 在 session 关闭前序列化)
@app.get("/users/{user_id}", response_model=UserPublic)
async def get_user(user_id: int, db: AsyncSession = Depends(get_db)):
    # get_db 中 session 在 yield 后才关闭,FastAPI 序列化在 yield 前完成
    return await db.get(User, user_id)

N+1 查询问题

# 陷阱:循环中触发多次查询
async def n_plus_one_problem(session: AsyncSession):
    users = await session.scalars(select(User))
    for user in users:
        print(user.posts)  # 每个 user 都触发一次额外查询!

# 正确做法:一次性预加载
async def no_n_plus_one(session: AsyncSession):
    stmt = select(User).options(selectinload(User.posts))
    users = await session.scalars(stmt)
    for user in users:
        print(user.posts)  # 已预加载,无额外查询

事务未提交

# 陷阱:忘记 commit
async def forgot_commit(session: AsyncSession):
    user = User(username="test")
    session.add(user)
    # 忘记 await session.commit()
    # session 关闭时自动 rollback,数据不会保存

# 使用 session.begin() 可以避免忘记提交
async def with_begin(session: AsyncSession):
    async with session.begin():
        user = User(username="test")
        session.add(user)
        # 自动 commit

autoflush 与查询一致性

# autoflush=True 时,执行 SELECT 前会自动 flush pending 的更改
# 这有时会导致意外行为

async def autoflush_example(session: AsyncSession):
    user = User(username="test", is_active=False)
    session.add(user)  # 还未 flush
    # 下一行 SELECT 触发 autoflush,user 被写入数据库(但未提交)
    count = await session.scalar(select(func.count(User.id)))
    # count 会包含 user,但事务未提交

最佳实践

selectinload / joinedload 消除 N+1 查询:访问关联对象前在查询中声明加载策略,而非让 SQLAlchemy 懒加载触发 N 次额外查询:

stmt = select(User).options(selectinload(User.posts))
users = (await session.scalars(stmt)).all()
# 2 次查询:先查 users,再 IN 查询 posts

Session 生命周期与请求绑定:在 FastAPI 中用 async with AsyncSession() as session 或依赖注入确保每个请求独占一个 Session,Session 不能跨请求或跨线程共享。

mapped_columnMapped 声明字段类型:SQLAlchemy 2.0 推荐用类型注解声明,IDE 支持更好:

class User(Base):
    __tablename__ = "users"
    id: Mapped[int] = mapped_column(primary_key=True)
    name: Mapped[str] = mapped_column(String(100))
    email: Mapped[str | None] = mapped_column(unique=True)

批量插入用 session.add_allinsert() 语句:单条循环 session.add 每次都刷新到数据库,大批量时性能差;用 await session.execute(insert(User), [dict, dict, ...]) 批量插入。

Alembic 迁移不要包含数据变更alembic revision --autogenerate 只生成 schema 变更,数据迁移写在独立脚本中,保证迁移可回滚且幂等。


常见陷阱

陷阱:在异步 Session 中混用同步查询

现象:async with AsyncSession() 中调用 session.execute(stmt)await,报 greenlet_spawn 或 coroutine 相关错误。
原因: AsyncSession 的所有 IO 操作都是协程,必须 await,不能像同步 Session 那样直接调用。
解决: 所有 session.execute()session.scalar()session.flush() 都加 await

陷阱:session.expire_on_commit=True 导致 commit 后访问属性报错

现象: await session.commit() 后访问 user.name 触发懒加载,在 async 环境下报 greenlet 错误或 MissingGreenlet
原因: 默认 expire_on_commit=True,commit 后所有 ORM 对象属性过期,访问时触发额外 SQL 但在异步上下文中无法同步执行。
解决: commit 前完成所有属性访问,或 session.expire_on_commit = False(但需要手动管理缓存一致性);更好的方案是 commit 后 await session.refresh(user) 重新加载需要的字段。

陷阱:back_populatesbackref 混用

现象: 在同一关系上同时设置 back_populatesbackref,启动时报 AmbiguousForeignKeysErrorInvalidRequestError
原因: backref 自动创建反向关系,若手动也设置了 back_populates,SQLAlchemy 认为关系重复定义。
解决: 只用一种:推荐 back_populates(显式声明双向关系,可读性好);不用 backref(隐式,难以追踪)。


参见

阅读更多

Web 安全基础

1. HTML 转义(服务端渲染必须): 2. CSP(Content Security Policy): 3. HttpOnly Cookie:防止 JS 读取会话 Cookie: 4. 前端框架防护: 攻击者在第三方网站构造一个表单,诱导已登录用户提交,浏览器会自动携带目标站的 Cookie。 触发条件: 1. 用户已登录目标网站(Cookie 有效) 2. 目标 API 仅凭 Cookie 识别用户身份 3. 请求来源未验证 1. CSRF Token(推荐): 2. SameSite Cookie: 3. 验证 Origin/Referer 头:

By yellowdog

HTTP 协议深度指南

HTTP(HyperText Transfer Protocol)是 Web 的基础传输协议,基于 TCP/IP,采用请求/响应模型。 相关文档:Web安全基础(/web-an-quan-ji-chu/) FastAPI完全指南(/fastapi-wan-quan-zhi-nan/) Nginx完全指南(/nginx-wan-quan-zhi-nan/) 幂等性:多次执行相同请求,服务器状态结果相同。PUT /users/1 多次执行结果一致;POST /users 每次创建新资源,非幂等。 浏览器直接从本地缓存读取,不向服务器发送请求。 缓存命中时,状

By yellowdog

系统设计基础

SLA 对照表: 选择建议:无状态服务(Web 层、API 层)优先水平扩展;数据库初期垂直扩展,达到瓶颈后考虑分库分表或读写分离。 缓存穿透(查询不存在的 key,每次都打到 DB): 缓存击穿(热点 key 过期,瞬间大量请求打到 DB): 缓存雪崩(大量 key 同时过期,或缓存服务宕机): 令牌桶 Python 实现: Redis 实现分布式限流(滑动窗口): URL 命名规则: Cursor 分页响应格式: 雪花算法结构(64 bit): 定义:分布式系统不能同时满足以下三个特性: 在分布式环境中 P 是必须保证的,所以实际是 CP vs AP

By yellowdog

算法思路与模板

二分查找要求序列有序,每次将搜索范围缩减一半,时间复杂度 O(log n)。 两个指针从两端向中间收缩,常用于有序数组。 滑动窗口维护一个满足条件的区间 left, right,right 不断向右扩张,条件不满足时收缩 left。 滑动窗口通用框架: 1. 确定"子问题":原问题可以分解为哪些规模更小的同类问题 2. 定义 dpi 或 dpij 的含义,要足够清晰 3. 推导状态转移方程 4. 确定初始状态(边界条件) 5. 确定计算顺序(确保依赖的子问题先计算) 每件物品最多选一次。dpj = 容量为 j 时的最大价值,逆序遍历容量防止重复选取。 每

By yellowdog