在标准 Django 项目中,Celery 并不能直接通过

python manage.py celery worker

启动,因为 Celery 官方推荐的是celery -A <proj> worker这种方式。

但在你的项目里,之所以可以用

python manage.py celery worker

启动,是因为项目集成了pipeline相关的 Celery 扩展(如pipeline.celery.settings),它为 Django注册了 celery 的相关命令。

celery启动常用参数

Celery Worker 常用参数:

基础参数

  • **A, --app** - 指定 Celery 应用实例
celery -A proj worker
  • **l, --loglevel** - 日志级别
-l info  # DEBUG, INFO, WARNING, ERROR, CRITICAL

队列相关

  • **Q, --queues** - 指定监听的队列(逗号分隔)会覆盖掉CELERY_QUEUES配置
-Q default,queue1,queue2
  • **X, --exclude-queues** - 排除某些队列
-X queue3,queue4

并发相关

  • **c, --concurrency** - 并发进程/线程数
-c 4  # 4个worker进程
  • **P, --pool** - 进程池实现方式
-P prefork  # prefork, eventlet, gevent, solo, threads
 
  • **-max-tasks-per-child** - 每个worker进程执行多少任务后重启
--max-tasks-per-child 1000

监控相关

  • **E, --task-events** - 开启任务事件监控
-E  # 用于 flower 或 celery events
  • **-heartbeat-interval** - 心跳间隔(秒)
--heartbeat-interval 30

其他常用

  • **n, --hostname** - worker名称
-n worker1@%h  # %h 是主机名
  • **-time-limit** - 任务硬超时(秒)
--time-limit 300  # 5分钟
  • **-soft-time-limit** - 任务软超时(秒)
--soft-time-limit 240  # 4分
  • **-pidfile** - PID文件路径
--pidfile /var/run/celery.pid
  • **-logfile** - 日志文件路径
--logfile /var/log/celery.log
  • **D, --detach** - 后台运行
-D  # daemon模式
  • **-autoscale** - 自动扩缩容
--autoscale=10,3  # 最大10个,最小3个

示例组合

# 完整示例
python manage.py celery worker \\
  -Q default,high_priority,low_priority \\
  -c 4 \\
  -P prefork \\
  -l info \\
  -n worker1@%h \\
  --time-limit=300 \\
  --max-tasks-per-child=1000 \\
  -E
 

查看所有参数

python manage.py celery worker --help
 

celery指定队列:

  1. 在被调用的task方法装饰器上
@task(queue='default')
def approve_callback(...):
    pass
  1. 调用时指定
approve_callback.apply_async(
    (auth.sn, status, provider_code, auth.ak, auth.sk, comment),
    queue='default'
)
  1. 在配置文件中配置
CELERY_ROUTES = {
    'approve.tasks.approve_callback': {'queue': 'default'},
    'approve.tasks.send_approve_mail': {'queue': 'default'},
}

一、准备工作:安装依赖

Celery 需依赖“消息代理(Broker)”和“结果存储(可选)”:

  • Broker:用于接收/分发任务(如 Redis、RabbitMQ,推荐 Redis 易配置);
  • 结果存储:用于保存任务执行结果(可选,如 Redis、数据库)。

1. 安装核心包

# 安装 Celery + Redis 依赖(Redis 作为 Broker 和结果存储)
pip install celery redis
# 若用 RabbitMQ 作为 Broker,需安装:pip install celery kombu==5.3.0 (适配版本)
 

2. 确保 Broker 服务可用

以 Redis 为例:

  • 本地安装 Redis 并启动(Windows 需手动下载,Linux 用 systemctl start redis);
  • 验证 Redis 连接:redis-cli ping,返回 PONG 表示正常。

二、项目配置:让 Celery 对接 Django

1. 核心配置文件(celery.py

在 Django 项目根目录(与 settings.py 同级)创建 celery.py,用于初始化 Celery 应用:

import os
from celery import Celery, platforms
from celery.signals import task_prerun, task_postrun
from django.conf import settings
from django.db import close_old_connections
 
# 1. 允许 root 用户启动 Worker(部署时可选,测试用)
platforms.C_FORCE_ROOT = True
 
# 2. 绑定 Django 配置模块(让 Celery 读取 Django 的 settings)
os.environ.setdefault("DJANGO_SETTINGS_MODULE", "your_project_name.settings")  # 替换为你的项目名
 
# 3. 初始化 Celery 应用(命名与项目一致)
app = Celery("your_project_name")
 
# 4. 从 Django settings 读取 Celery 配置(仅读取 CELERY_ 前缀的参数)
app.config_from_object("django.conf:settings", namespace="CELERY")
 
# 5. 自动发现所有 Django App 中的任务(任务定义在 app/tasks.py 中)
app.autodiscover_tasks(lambda: settings.INSTALLED_APPS)
 
# 6. 任务前后关闭过期数据库连接(避免连接泄漏)
@task_prerun.connect
@task_postrun.connect
def close_db_connections(**kwargs):
    close_old_connections()
 

2. 配置 Django settings.py

settings.py 中添加 Celery 核心参数(Broker、结果存储、时区等):

# ---------------------- Celery 配置 ----------------------
# 1. Broker 地址(Redis 示例,格式:redis://[密码@]主机:端口/数据库号)
CELERY_BROKER_URL = "redis://localhost:6379/0"  # 数据库号 0(避免与其他服务冲突)
 
# 2. 任务结果存储地址(与 Broker 共用 Redis,可选)
CELERY_RESULT_BACKEND = "redis://localhost:6379/0"
 
# 3. 时区(必须与 Django 时区一致,避免定时任务时间偏差)
CELERY_TIMEZONE = "Asia/Shanghai"
 
# 4. 任务结果过期时间(单位:秒,避免 Redis 存储膨胀)
CELERY_RESULT_EXPIRES = 3600  # 1 小时后自动删除结果
 
# 5. Worker 启动时重试 Broker 连接(Celery 6.0+ 新增,提高稳定性)
CELERY_BROKER_CONNECTION_RETRY_ON_STARTUP = True
 

3. 项目入口文件(__init__.py

在项目根目录的 __init__.py 中导入 Celery 应用,确保 Django 启动时能加载:

# your_project_name/__init__.py
from .celery import app as celery_app
 
__all__ = ("celery_app",)
 

三、开发任务:定义/调用异步任务

1. 定义任务(app/tasks.py

在任意 Django App 下创建 tasks.py(Celery 会自动扫描该文件),定义异步任务:

# 示例:在 blog 应用下定义任务
from celery import shared_task
from django.core.mail import send_mail
from your_project_name.settings import DEFAULT_FROM_EMAIL
 
# 用 @shared_task 装饰器定义“可被所有 App 调用的任务”
@shared_task(bind=True, retry_kwargs={"max_retries": 3})  # 绑定任务实例,允许重试(最多 3 次)
def send_notify_email(self, to_email, subject, message):
    """异步发送通知邮件(耗时操作,避免阻塞请求)"""
    try:
        send_mail(
            subject=subject,
            message=message,
            from_email=DEFAULT_FROM_EMAIL,
            recipient_list=[to_email],
            fail_silently=False,
        )
        return f"邮件已发送至 {to_email}"  # 任务成功返回结果
    except Exception as e:
        # 捕获异常时重试(每次重试间隔 5 秒)
        self.retry(exc=e, countdown=5)
 
  • **@shared_task**:Celery 用于定义“跨 App 可调用任务”的装饰器(无需绑定特定 Celery 应用);
  • **bind=True**:将任务实例绑定到第一个参数 self,可调用 self.retry() 实现重试;
  • 任务逻辑:封装耗时/异步操作(如发送邮件、生成报表、调用第三方接口)。

2. 调用任务(在视图/服务中)

在 Django 视图、模型方法或服务中,通过 delay()apply_async() 调用异步任务:

# 示例:在视图中调用发送邮件任务
from django.http import JsonResponse
from blog.tasks import send_notify_email
 
def user_register(request):
    # 1. 处理用户注册逻辑(同步操作,快速完成)
    username = request.POST.get("username")
    email = request.POST.get("email")
    # ... 保存用户到数据库 ...
 
    # 2. 异步调用发送邮件任务(不阻塞当前请求)
    # delay():简单调用,参数直接传递
    task_result = send_notify_email.delay(
        to_email=email,
        subject="注册成功",
        message=f"欢迎 {username} 注册!"
    )
 
    # 3. 返回响应(任务已提交到 Broker,正在排队执行)
    return JsonResponse({
        "code": 200,
        "message": "注册成功,邮件正在发送",
        "task_id": task_result.id  # 任务唯一 ID,可用于查询结果
    })
 
  • **delay(*args, **kwargs)**:简化版调用,直接传递任务参数;
  • **apply_async(args, kwargs, countdown=10)**:高级调用,支持延迟执行(countdown=10 表示 10 秒后执行)、指定队列等。

四、运行与管理:启动 Worker 执行任务

Celery 任务需通过“Worker 进程”执行——Worker 会持续监听 Broker 中的任务队列,拿到任务后立即执行。

1. 启动 Worker(终端命令)

在项目根目录(manage.py 同级)执行:

# 基础启动:启动 1 个 Worker,监听所有队列
celery -A your_project_name worker --loglevel=info
 
# 常用参数扩展:
# 1. 启动 4 个 Worker 进程(提高并发,适合多核服务器)
celery -A your_project_name worker -c 4 --loglevel=info
 
# 2. 指定队列(仅执行特定队列的任务,适合任务分类)
celery -A your_project_name worker -Q email_queue --loglevel=info
 
# 3. 后台运行(Linux 用 nohup,日志输出到 celery.log)
nohup celery -A your_project_name worker -c 4 --loglevel=info > celery.log 2>&1 &
 
  • **A your_project_name**:指定 Celery 应用(对应 celery.py 中的初始化应用);
  • **-loglevel=info**:日志级别(info 显示任务执行状态,debug 显示详细信息);
  • **c 4**:Worker 进程数(建议与 CPU 核心数一致,避免资源浪费)。

2. 查看任务结果(可选)

若配置了 CELERY_RESULT_BACKEND,可通过任务 ID 查询结果:

# 在 Django  shell 中查询
python manage.py shell
 
from celery.result import AsyncResult
from your_project_name.celery import app
 
# 替换为实际任务 ID(如 user_register 视图返回的 task_id)
task_id = "a1b2c3d4-5678-90ef-ghij-klmnopqrstuv"
result = AsyncResult(task_id, app=app)
 
print(result.state)  # 任务状态:PENDING(排队)、SUCCESS(成功)、FAILURE(失败)
if result.state == "SUCCESS":
    print(result.result)  # 任务成功返回的结果(如“邮件已发送至 xxx@xxx.com”)
elif result.state == "FAILURE":
    print(result.traceback)  # 任务失败的异常栈信息
 

3. 定时任务(可选:Celery Beat)

若需要“定时执行任务”(如每天凌晨同步数据),需配置 Celery Beat(定时任务调度器):

  1. settings.py 中添加定时任务:
# Celery Beat 定时任务配置
CELERY_BEAT_SCHEDULE = {
    "daily-data-sync": {  # 任务名称(自定义)
        "task": "blog.tasks.sync_daily_data",  # 要执行的任务(路径+函数名)
        "schedule": 86400,  # 执行间隔(单位:秒,86400 表示 1 天)
        # "schedule": crontab(hour=3, minute=0),  # 用 crontab 表达式(每天凌晨 3 点执行)
        "args": ("2024-01-01",),  # 任务参数
    },
}
 
  1. 启动 Beat + Worker(需同时启动,Beat 负责发任务,Worker 负责执行):
celery -A your_project_name beat -l info  # 单独启动 Beat(仅调度,不执行)
celery -A your_project_name worker -B -l info  # 同时启动 Beat 和 Worker(测试用)
 

五、常见问题与注意事项

  1. Worker 启动报错“Connection refused”:检查 Broker(如 Redis)是否启动,地址是否正确;
  2. 任务执行时报“数据库连接错误”:确保 celery.py 中配置了 task_prerun/task_postrun 关闭过期连接;
  3. 任务不执行:检查 Worker 是否启动、是否监听了正确的队列,Broker 中是否有任务堆积(redis-cli KEYS "*celery*" 查看);
  4. 生产环境部署:用 supervisorsystemd 管理 Worker 进程(避免终端关闭后 Worker 退出),定期清理过期任务结果。

总结

Django + Celery 核心流程可概括为:

配置联动(celery.py + settings)定义任务(tasks.py)调用任务(视图/服务)启动 Worker 执行任务

核心价值是“将耗时/异步操作剥离出 Django 主进程”,避免阻塞 HTTP 请求,提升项目响应速度和稳定性。

相关文档