FastAPI项目中Celery与RabbitMQ队列的合理配置原则
在FastAPI+SQLAlchemy项目中集成Celery的实践方案
一、Celery基础配置(适配现有项目结构)
你的项目已有celery目录,直接在该目录下完成配置:
- 创建配置文件:
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, } }
- 初始化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
五、核心注意事项
- 任务幂等性:确保任务重复执行不会产生副作用(比如批量导入任务需校验数据是否已存在)
- 异常处理与重试:对网络依赖(比如邮件服务)、数据库操作添加异常捕获,合理配置重试策略
- 任务状态查询:若需前端查询任务状态,可通过
celery_app.AsyncResult(task_id)获取状态和结果,需在FastAPI中添加对应的查询路由 - 数据库连接池:确保SQLAlchemy的连接池配置适配Celery Worker的并发数,避免连接耗尽
内容的提问来源于stack exchange,提问作者Genry
相关产品推荐
相关产品推荐

