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/