Python FastAPI异步接口实战:async/await高并发处理与Pydantic数据校验

FastAPI基于Starlette ASGI框架和Pydantic数据模型构建,原生支持async/await异步IO,单进程即可处理数千并发请求。与Flask的同步模型相比,在IO密集型场景下QPS提升数倍。本文覆盖FastAPI的核心能力:异步路由处理、Pydantic模型校验、依赖注入、后台任务和高并发部署。

Pydantic模型定义与请求数据自动校验

FastAPI使用Pydantic BaseModel声明接口的数据结构,框架自动完成请求体解析、类型校验和文档生成。定义模型时支持嵌套、枚举、约束和自定义校验器:

from pydantic import BaseModel, Field, field_validator, EmailStr
from enum import Enum
from datetime import datetime
from typing import Optional

class OrderStatus(str, Enum):
    pending = "pending"
    paid = "paid"
    shipped = "shipped"
    cancelled = "cancelled"

class OrderItem(BaseModel):
    product_id: int = Field(gt=0, description="商品ID")
    quantity: int = Field(ge=1, le=99, description="数量")
    unit_price: float = Field(gt=0, description="单价")

    @field_validator('quantity')
    @classmethod
    def check_quantity(cls, v):
        if v > 50:
            raise ValueError('单笔数量不能超过50')
        return v

class CreateOrderRequest(BaseModel):
    user_email: EmailStr
    items: list[OrderItem] = Field(min_items=1, max_items=20)
    shipping_address: str = Field(min_length=10, max_length=500)
    remark: Optional[str] = Field(None, max_length=200)

class OrderResponse(BaseModel):
    order_id: str
    status: OrderStatus
    total_amount: float
    created_at: datetime
    items: list[OrderItem]

路由函数直接使用模型作为类型注解,FastAPI自动校验请求体并返回422错误:

from fastapi import FastAPI, HTTPException

app = FastAPI(title="订单服务", version="1.0.0")

@app.post("/orders", response_model=OrderResponse, status_code=201)
async def create_order(req: CreateOrderRequest):
    total = sum(item.quantity * item.unit_price for item in req.items)
    order = OrderResponse(
        order_id="ORD-" + datetime.now().strftime("%Y%m%d%H%M%S"),
        status=OrderStatus.pending,
        total_amount=round(total, 2),
        created_at=datetime.now(),
        items=req.items,
    )
    # 持久化逻辑...
    return order

非法请求自动拒绝,如quantity=-1返回{“type”:”greater_than_equal”,”loc”:[“body”,”items”,0,”quantity”],”msg”:”Input should be greater than or equal to 1″}。校验逻辑与业务逻辑解耦,减少手写if-else校验代码。

async/await异步数据库操作与连接池配置

同步数据库驱动会阻塞事件循环,抵消async优势。使用SQLAlchemy 2.0的async引擎配合asyncpg(PostgreSQL)或aiomysql:

from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import sessionmaker, DeclarativeBase

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

engine = create_async_engine(
    DATABASE_URL,
    pool_size=20,
    max_overflow=10,
    pool_pre_ping=True,
    pool_recycle=3600,
)
AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)

class Base(DeclarativeBase):
    pass

# 依赖注入:每个请求获取独立session
from fastapi import Depends

async def get_db():
    async with AsyncSessionLocal() as session:
        try:
            yield session
            await session.commit()
        except Exception:
            await session.rollback()
            raise
        finally:
            await session.close()

@app.get("/orders/{order_id}", response_model=OrderResponse)
async def get_order(order_id: str, db: AsyncSession = Depends(get_db)):
    result = await db.execute(
        text("SELECT * FROM orders WHERE order_id = :oid"),
        {"oid": order_id}
    )
    row = result.fetchone()
    if not row:
        raise HTTPException(status_code=404, detail="订单不存在")
    return OrderResponse(**dict(row._mapping))

pool_size=20设置连接池常驻连接数,max_overflow=10允许短时超出到30个连接。pool_pre_ping在取连接前做健康检查,避免使用已断开的连接。pool_recycle=3600定期回收连接,防止数据库端超时关闭。

Background Tasks与异步任务队列选型

轻量后台任务用FastAPI内置的BackgroundTasks,请求返回后异步执行:

from fastapi import BackgroundTasks
import smtplib
from email.mime.text import MIMEText

async def send_notification(email: str, order_id: str):
    msg = MIMEText(f"您的订单 {order_id} 已创建")
    msg['Subject'] = '订单创建通知'
    msg['To'] = email
    # 异步发送邮件
    with smtplib.SMTP('smtp.example.com', 587) as smtp:
        smtp.send_message(msg)

@app.post("/orders", response_model=OrderResponse, status_code=201)
async def create_order(
    req: CreateOrderRequest,
    background_tasks: BackgroundTasks,
    db: AsyncSession = Depends(get_db),
):
    # ... 创建订单逻辑 ...
    background_tasks.add_task(send_notification, req.user_email, order.order_id)
    return order

BackgroundTasks适合耗时短(秒级)且非关键的任务。对于可靠性要求高的任务(如订单支付回调、批量数据处理),应使用专业任务队列。Celery+Redis方案成熟但较重,arq(基于Redis的轻量async任务队列)与FastAPI的async生态更契合:

# tasks.py
from arq import create_pool
from arq.connections import RedisSettings

async def process_payment(ctx, order_id: str):
    # 耗时支付处理
    ...

class WorkerSettings:
    functions = [process_payment]
    redis_settings = RedisSettings(host='localhost', port=6379)
    max_jobs = 10
    job_timeout = 300

# 路由中投递任务
from arq import create_pool

@app.post("/orders/{order_id}/pay")
async def pay_order(order_id: str):
    redis = await create_pool(RedisSettings())
    await redis.enqueue_job('process_payment', order_id)
    return {"status": "processing"}

Gunicorn+Uvicorn生产部署与性能压测

开发环境用uvicorn直接启动,生产环境推荐Gunicorn管理多个Uvicorn worker进程:

gunicorn main:app \
  -w 4 \
  -k uvicorn.workers.UvicornWorker \
  --bind 0.0.0.0:8000 \
  --access-logfile - \
  --error-logfile - \
  --graceful-timeout 30 \
  --timeout 60

-w参数设置worker进程数,建议设为CPU核心数的2-4倍。每个worker内部为单进程异步事件循环,worker间通过Gunicorn的预派发模型隔离。–graceful-timeout确保优雅关闭时正在处理的请求能完成。

使用wrk做压测对比同步框架差异:

# FastAPI async endpoint
wrk -t4 -c100 -d30s http://localhost:8000/orders/ORD-001

# 典型结果:
# Requests/sec:  12000+
# Latency:       8.2ms (avg)
# Transfer/sec:  2.1MB

CPU密集型任务需用run_in_executor offload到线程池,避免阻塞事件循环:

import asyncio
from concurrent.futures import ThreadPoolExecutor

executor = ThreadPoolExecutor(max_workers=4)

@app.post("/export")
async def export_report():
    loop = asyncio.get_event_loop()
    result = await loop.run_in_executor(
        executor, 
        cpu_heavy_excel_generation  # 同步CPU密集函数
    )
    return {"url": result}

worker数量、连接池大小和线程池大小需要配合调优。一般原则:async IO并发数远大于worker数,worker数匹配CPU核数,线程池仅用于不可避免的同步阻塞调用。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/pythonfastapi-yi-bu-jie-kou-shi-zhan-asyncawait-gao-bing-fa/

(0)
小编小编
上一篇 2小时前
下一篇 2小时前

相关推荐