统一采用
session: AsyncSession = Depends(get_session)依赖注入方式,不使用中间件request.state.session;事务使用async with session.begin(),自动 com
依赖准备(前面已写)
from fastapi import FastAPI, Depends, HTTPException from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import select, delete, update from pydantic import BaseModel from typing import List from models import AsyncSessionFactory, User app = FastAPI() # 依赖函数:获取数据库Session async def get_session(): session = AsyncSessionFactory() try: yield session finally: await session.close() # 响应Schema class UserRespSchema(BaseModel): id: int email: str username: str class Config: orm_mode = True # 创建用户请求Schema class UserCreateReqSchema(BaseModel): email: str username: str password: str1、新增数据 create
@app.post("/user/add", response_model=UserRespSchema) async def add_user( req: UserCreateReqSchema, session: AsyncSession = Depends(get_session) ): try: async with session.begin(): user = User(username=req.username, email=req.email, password=req.password) session.add(user) return user except Exception: raise HTTPException(status_code=400, detail="用户名或邮箱已经存在!")2、删除数据 delete
@app.delete("/user/delete/{user_id}") async def delete_user( user_id: int, session: AsyncSession = Depends(get_session) ): async with session.begin(): await session.execute( delete(User).where(User.id == user_id) ) return {"message": "删除成功!"}3、查询数据 select
查询单条
@app.get("/user/select/{user_id}", response_model=UserRespSchema) async def select_user_by_id( user_id: int, session: AsyncSession = Depends(get_session) ): async with session.begin(): query = await session.execute( select(User.id, User.email, User.username).where(User.id == user_id) ) result = query.one()._asdict() return result查询多条 + 模糊搜索、分页排序
@app.get("/user/list", response_model=List[UserRespSchema]) async def select_user_list( q: str | None = None, session: AsyncSession = Depends(get_session) ): from sqlalchemy import or_ stmt = select(User.id, User.username, User.email) if q: stmt = stmt.where( or_( User.email.contains(q), User.username.contains(q) ) ) # 分页、排序 stmt = stmt.limit(2).offset(0).order_by(User.id.desc()) query = await session.execute(stmt) result = [row._asdict() for row in query] return result注意:
session.query()只适用于同步 SQLAlchemy,select()。
4、修改数据 update
方式 1:先查询对象,修改对象属性(执行 2 次 SQL,可以直接返回修改后对象)
@app.put("/user/update/{user_id}", response_model=UserRespSchema) async def update_user_obj( user_id: int, user_data: UserCreateReqSchema, session: AsyncSession = Depends(get_session) ): async with session.begin(): query = await session.execute(select(User).where(User.id == user_id)) user = query.scalars().first() user.email = user_data.email user.username = user_data.username return user方式 2:直接 update 语句(只执行 1 次 SQL,不会刷新 ORM 对象)
@app.put("/user/update_direct/{user_id}") async def update_user_direct( user_id: int, user_data: UserCreateReqSchema, session: AsyncSession = Depends(get_session) ): async with session.begin(): await session.execute( update(User) .where(User.id == user_id) .values(**user_data.dict()) ) return {"message": "数据修改成功!"}