You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

FastAPI项目中Celery与RabbitMQ队列的合理配置原则

在FastAPI+SQLAlchemy项目中集成Celery的实践方案

一、Celery基础配置(适配现有项目结构)

你的项目已有celery目录,直接在该目录下完成配置:

  1. 创建配置文件:celery/config.py
# 复用项目已有的Redis配置
REDIS_URL = "redis://localhost:6379/0"

CELERY_CONFIG = {
    "broker_url": REDIS_URL,
    "result_backend": REDIS_URL,
    # 任务序列化格式,推荐json(避免pickle安全问题)
    "task_serializer": "json",
    "result_serializer": "json",
    "accept_content": ["json"],
    # 时区和项目保持一致
    "timezone": "Asia/Shanghai",
    "enable_utc": False,
    # 任务过期时间,避免死任务堆积
    "task_expires": 3600,
    # 自动重试配置(可选)
    "task_retry_policy": {
        "max_retries": 3,
        "interval_start": 0,
        "interval_step": 0.5,
        "interval_max": 2,
    }
}
  1. 初始化Celery实例:celery/main.py
from celery import Celery
from celery.config import from_object
from celery.schedules import crontab

from .config import CELERY_CONFIG

# 初始化Celery,指定项目模块路径
celery_app = Celery("company_employee_management")
celery_app.config_from_object(CELERY_CONFIG)

# 自动发现tasks目录下的任务
celery_app.autodiscover_tasks(["celery.tasks"])

# 可选:配置定时任务(比如每日统计公司员工数量)
celery_app.conf.beat_schedule = {
    "daily-company-employee-stats": {
        "task": "celery.tasks.company_tasks.daily_employee_stats",
        "schedule": crontab(hour=0, minute=0),
    },
}

二、任务定义与目录规范

按业务模块拆分任务,在celery目录下新建tasks子目录,对应你的公司、员工管理场景:

  • celery/tasks/company_tasks.py:处理公司相关异步任务
  • celery/tasks/employee_tasks.py:处理员工相关异步任务
  • celery/tasks/notification_tasks.py:处理通知类异步任务

任务编写注意事项:

任务内的数据库操作需独立创建SQLAlchemy会话,不能复用FastAPI进程的会话(Celery Worker是独立进程)

# 示例:employee_tasks.py中的入职通知任务
from sqlalchemy.orm import Session
from celery.main import celery_app

from db.session import get_db  # 复用项目的DB会话创建函数
from models.employee import Employee
from models.company import Company
from services.notification import send_email  # 你的通知服务

@celery_app.task(bind=True, name="send-employee-onboarding-notification")
def send_employee_onboarding_notification(self, employee_id: int, company_id: int):
    db: Session = next(get_db())
    try:
        employee = db.query(Employee).filter(Employee.id == employee_id).first()
        company = db.query(Company).filter(Company.id == company_id).first()
        if not employee or not company:
            raise ValueError("员工或公司不存在")
        
        # 执行异步通知逻辑(比如发送入职邮件)
        send_email(
            to=employee.email,
            subject=f"欢迎加入{company.name}",
            content=f"您好{employee.name},欢迎成为{company.name}的一员..."
        )
        return {"status": "success", "message": "通知发送完成"}
    except Exception as e:
        # 触发重试(根据配置的重试策略)
        self.retry(exc=e)
    finally:
        db.close()

三、FastAPI与Celery的集成

在main.py或对应业务路由中直接调用Celery任务,注意:路由层已完成JWT权限校验,任务无需重复校验(若任务需用户信息,可将必要参数传入任务)

示例:员工创建路由中触发异步通知

from fastapi import APIRouter, Depends, status
from sqlalchemy.orm import Session
from auth.jwt import get_current_user  # 你的JWT依赖

from db.session import get_db
from schemas.employee import EmployeeCreate
from services.employee import create_employee
from celery.tasks.employee_tasks import send_employee_onboarding_notification

router = APIRouter(prefix="/employees", tags=["Employees"])

@router.post("/", status_code=status.HTTP_201_CREATED)
def create_new_employee(
    employee_data: EmployeeCreate,
    db: Session = Depends(get_db),
    current_user = Depends(get_current_user)
):
    # 路由层已做权限校验,确保当前用户有创建员工权限
    new_employee = create_employee(db=db, employee_data=employee_data, company_id=current_user.company_id)
    
    # 触发异步任务,使用delay()方法
    task = send_employee_onboarding_notification.delay(new_employee.id, current_user.company_id)
    
    # 返回任务ID,供前端查询任务状态(可选)
    return {"employee_id": new_employee.id, "task_id": task.id, "status": "employee created, notification pending"}

四、队列划分逻辑

根据你的业务场景和角色权限,按业务类型+优先级划分队列,避免不同任务互相阻塞:

推荐队列划分:

  • default队列:处理通用、低优先级的零散任务(比如单个员工信息更新后的同步操作)
  • company_ops队列:处理公司相关的批量/耗时任务(比如批量创建子公司、公司数据归档、月度报表生成)
  • employee_ops队列:处理员工相关的批量/耗时任务(比如批量导入员工、薪资计算、员工数据批量导出)
  • notification队列:处理所有通知类任务(邮件、短信通知)

任务绑定队列:

定义任务时指定队列,或在配置中设置路由规则:

# 方式1:定义任务时直接指定队列
@celery_app.task(bind=True, name="batch-import-employees", queue="employee_ops")
def batch_import_employees(self, file_path: str, company_id: int):
    # 批量导入逻辑
    pass

# 方式2:配置路由规则(在celery/config.py的CELERY_CONFIG中添加)
CELERY_CONFIG.update({
    "task_routes": {
        "celery.tasks.company_tasks.*": {"queue": "company_ops"},
        "celery.tasks.employee_tasks.batch_import*": {"queue": "employee_ops"},
        "celery.tasks.notification_tasks.*": {"queue": "notification"},
    }
})

启动Worker时指定监听队列:

# 监听所有队列
celery -A celery.main worker -l info -Q default,company_ops,employee_ops,notification

# 仅监听员工操作队列(可单独部署Worker处理高负载任务)
celery -A celery.main worker -l info -Q employee_ops --concurrency=4

五、核心注意事项

  1. 任务幂等性:确保任务重复执行不会产生副作用(比如批量导入任务需校验数据是否已存在)
  2. 异常处理与重试:对网络依赖(比如邮件服务)、数据库操作添加异常捕获,合理配置重试策略
  3. 任务状态查询:若需前端查询任务状态,可通过celery_app.AsyncResult(task_id)获取状态和结果,需在FastAPI中添加对应的查询路由
  4. 数据库连接池:确保SQLAlchemy的连接池配置适配Celery Worker的并发数,避免连接耗尽

内容的提问来源于stack exchange,提问作者Genry

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 07:07:48