下载工作台
FastAPI 开发

BackgroundTasks 与 Celery 异步

试读上半部分 · 解锁后可读全文

第 10 章 · BackgroundTasks 与 Celery 异步

本章目标:对比 FastAPI BackgroundTasksCelery 分布式任务队列;配置 Redis 作为 Celery broker 与 result backend;实现「下单/上传后异步通知」示例;理解何时用进程内后台任务、何时上 Celery Worker;掌握任务重试、超时与可观测基础。

学时建议:5~6 小时(含 2 小时 Celery + Redis 跟练)

前置:完成 fastapi-web ch09(上传接口);ch07 用户体系;本机可运行 Redis(Docker 或本地安装)。


10.1 同步 API 的耗时陷阱

以下操作若放在请求线程内同步执行,会拉长响应时间:

操作典型耗时是否适合阻塞请求
写数据库5~50 ms通常保留在请求内
发送邮件/短信0.5~5 s应异步
生成缩略图100 ms~2 s可异步
调用第三方 webhook不确定应异步
批量导出 CSV数秒~数分钟必须异步
同步(差体验)                    异步(推荐)
POST /orders  ──► 支付 ──► 发邮件 ──► 200     POST /orders ──► 支付 ──► 200
                      (用户等 3 秒)              └──► 队列 ──► Worker 发邮件

svc-demo 场景:商品封面上传成功后 发送运营通知;用户注册后 发送欢迎邮件(模拟)。


10.2 BackgroundTasks:轻量进程内异步

FastAPI 内置,同一进程、请求返回后执行,无持久化、无分布式

app/services/notifications.py

import logging

logger = logging.getLogger("svc.notify")


def send_welcome_email(email: str, name: str | None) -> None:
    # 模拟 SMTP;生产对接邮件服务
    logger.info("welcome email sent to %s (%s)", email, name or "-")


def notify_cover_updated(product_id: int, cover_url: str) -> None:
    logger.info("ops notified: product=%s cover=%s", product_id, cover_url)

app/api/v1/auth.py 注册后:

from fastapi import BackgroundTasks
from app.services.notifications import send_welcome_email


@router.post("/register", response_model=UserRead, status_code=201)
async def register(
    body: UserCreate,
    background_tasks: BackgroundTasks,
    repo: UserRepository = Depends(get_user_repo),
):
    if await repo.get_by_email(body.email):
        raise ConflictError(message="邮箱已注册", code="DUPLICATE_EMAIL")
    user = await repo.create(body.email, body.password, body.full_name)
    background_tasks.add_task(send_welcome_email, user.email, user.full_name)
    return user

上传封面(ch09):

@router.post("/{product_id}/cover", response_model=ProductRead)
async def upload_product_cover(
    ...,
    background_tasks: BackgroundTasks,
):
    # ... 保存图片 ...
    background_tasks.add_task(notify_cover_updated, product.id, cover_url)
    return await repo.update(product, cover_url=cover_url)
优点缺点
零依赖、代码简单进程重启任务丢失
适合 <1s 收尾工作不能水平扩展 Worker
与 async 路由兼容长任务阻塞 worker 进程

适用:写审计日志、轻量通知、缓存失效——可丢、可短、单实例


10.3 Celery:分布式任务队列

FastAPI (Producer)          Redis (Broker)           Celery Worker
    .delay(task)    ──►    队列 list/stream    ──►   消费执行
                                │
                                └──► result backend(可选)
组件作用
Broker消息中间件,本章用 Redis
Worker独立进程执行任务
Backend存任务结果(可用 Redis 或禁用)
Beat定时任务调度(扩展)

安装:

pip install "celery[redis]>=5.3,<6"
pip install "redis>=5,<6"

Docker 启动 Redis(练习环境):

docker run -d --name svc-redis -p 6379:6379 redis:7-alpine

10.4 Celery 应用配置

项目根目录 celery_app.py

from celery import Celery

from app.core.config import settings

celery_app = Celery(
    "svc_demo",
    broker=settings.CELERY_BROKER_URL,
    backend=settings.CELERY_RESULT_BACKEND,
    include=["app.tasks.notifications"],
)

celery_app.conf.update(
    task_serializer="json",
    accept_content=["json"],
    result_serializer="json",
    timezone="Asia/Shanghai",
    enable_utc=True,
    task_track_started=True,
    task_time_limit=300,
    task_soft_time_limit=240,
    broker_connection_retry_on_startup=True,
)

app/core/config.py

class Settings(BaseSettings):
    # ...
    CELERY_BROKER_URL: str = "redis://127.0.0.1:6379/0"
    CELERY_RESULT_BACKEND: str = "redis://127.0.0.1:6379/1"
    NOTIFY_WEBHOOK_URL: str = "https://hooks.example.com/ops-alerts"

10.5 定义 Celery 任务

app/tasks/notifications.py

import logging
from typing import Any

import httpx
from celery_app import celery_app

logger = logging.getLogger("svc.celery")


@celery_app.task(bind=True, max_retries=3, default_retry_delay=10)
def send_ops_webhook(self, payload: dict[str, Any]) -> dict:
    """异步通知运营 webhook(虚构地址)。"""
    from app.core.config import settings

    url = settings.NOTIFY_WEBHOOK_URL
    try:
        with httpx.Client(timeout=10.0) as client:
            resp = client.post(url, json=payload)
            resp.raise_for_status()
        logger.info("webhook ok product=%s", payload.get("product_id"))
        return {"status": "ok", "code": resp.status_code}
    except Exception as exc:
        logger.warning("webhook failed attempt=%s: %s", self.request.retries, exc)
        raise self.retry(exc=exc)


@celery_app.task
def generate_thumbnail(product_id: int, source_path: str) -> str:

以下内容需解锁后阅读

试读已结束。解锁本章 ¥5.00,或开通年度会员畅读全部教程。
年度会员 ¥199.00/年; 小紫 AI 工作台有效会员 ¥99.00/年

正文仅在服务端鉴权后下发,未付费无法获取下半部分内容。