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_column 和 Mapped 声明字段类型: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_all 或 insert() 语句:单条循环 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_populates 和 backref 混用
现象: 在同一关系上同时设置 back_populates 和 backref,启动时报 AmbiguousForeignKeysError 或 InvalidRequestError。
原因: backref 自动创建反向关系,若手动也设置了 back_populates,SQLAlchemy 认为关系重复定义。
解决: 只用一种:推荐 back_populates(显式声明双向关系,可读性好);不用 backref(隐式,难以追踪)。