django-celery

django-celery

热门

Django + Celery 异步任务开发最佳实践——涵盖项目配置、任务设计、Beat 定时调度、重试机制、Canvas 工作流、服务监控及单元测试。适用于在 Django 项目中引入后台任务、定时任务或异步处理等场景。

24万Star
3.6万Fork
更新于 2026/8/2
SKILL.md
只读
名称
django-celery
描述

Django + Celery 异步任务开发最佳实践——涵盖项目配置、任务设计、Beat 定时调度、重试机制、Canvas 工作流、服务监控及单元测试。适用于在 Django 项目中引入后台任务、定时任务或异步处理等场景。

Django + Celery 异步任务最佳实践模式

使用 Celery 配合 Redis 或 RabbitMQ 在 Django 中构建生产级后台任务处理架构的模式与实践。

触发条件 / 使用场景

  • 为 Django 应用添加后台任务或异步处理能力
  • 实现周期性任务或定时调度任务
  • 将耗时操作(如发送邮件、生成 PDF、调用第三方 API)从请求响应循环中剥离解耦
  • 配置 Celery Beat 以实现类 cron 的定时任务调度
  • 排查任务失败、重试机制或消息队列积压问题
  • 为 Celery 任务编写单元测试与集成测试

项目配置

环境安装

pip install 'celery[redis]' django-celery-results django-celery-beat

celery.py — 应用入口点

# config/celery.py
import os
from celery import Celery

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'config.settings.development')

app = Celery('myproject')
app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks()  # 自动发现各个 INSTALLED_APP 中的 tasks.py

@app.task(bind=True, ignore_result=True)
def debug_task(self):
    print(f'Request: {self.request!r}')
# config/__init__.py
from .celery import app as celery_app

__all__ = ('celery_app',)

Django 配置项

# config/settings/base.py

# Broker 消息中间件(生产环境推荐使用 Redis)
CELERY_BROKER_URL = env('CELERY_BROKER_URL', default='redis://localhost:6379/0')
CELERY_RESULT_BACKEND = env('CELERY_RESULT_BACKEND', default='django-db')

# 序列化设置
CELERY_ACCEPT_CONTENT = ['json']
CELERY_TASK_SERIALIZER = 'json'
CELERY_RESULT_SERIALIZER = 'json'

# 任务行为设置
CELERY_TASK_TRACK_STARTED = True
CELERY_TASK_TIME_LIMIT = 30 * 60        # 硬性超时限制:30 分钟
CELERY_TASK_SOFT_TIME_LIMIT = 25 * 60   # 软性超时限制:达到 25 分钟时抛出 SoftTimeLimitExceeded
CELERY_WORKER_PREFETCH_MULTIPLIER = 1   # 防止 Worker 预取过多长耗时任务导致积压
CELERY_TASK_ACKS_LATE = True            # Worker 崩溃时将任务重新放回队列

# 结果持久化
CELERY_RESULT_EXPIRES = 60 * 60 * 24   # 结果保留 24 小时

# Beat 调度器(用于定时任务)
CELERY_BEAT_SCHEDULER = 'django_celery_beat.schedulers:DatabaseScheduler'

# 已安装的应用
INSTALLED_APPS += [
    'django_celery_results',
    'django_celery_beat',
]

运行 Worker 服务

# 启动 Worker(开发环境)
celery -A config worker --loglevel=info

# 启动 Beat 调度器(定时任务)
celery -A config beat --loglevel=info --scheduler django_celery_beat.schedulers:DatabaseScheduler

# 混合启动 Worker 与 Beat(仅限开发环境,切勿用于生产环境)
celery -A config worker --beat --loglevel=info

# 生产环境:多 Worker 并发并指定优先级队列
celery -A config worker --loglevel=warning --concurrency=4 -Q default,high_priority

任务设计模式

基础任务

# apps/notifications/tasks.py
from celery import shared_task
import logging

logger = logging.getLogger(__name__)

@shared_task(name='notifications.send_welcome_email')
def send_welcome_email(user_id: int) -> None:
    """向新注册用户发送欢迎邮件。"""
    from apps.users.models import User
    from apps.notifications.services import EmailService

    try:
        user = User.objects.get(pk=user_id)
    except User.DoesNotExist:
        logger.warning('send_welcome_email: user %s not found', user_id)
        return  # 保证幂等性——直接返回不抛错,此时任务已无法执行

    EmailService.send_welcome(user)
    logger.info('已向用户 %s 发送欢迎邮件', user_id)

可重试任务

@shared_task(
    bind=True,
    name='integrations.sync_to_crm',
    max_retries=5,
    default_retry_delay=60,       # 首次重试前的等待秒数
    autoretry_for=(ConnectionError, TimeoutError),
    retry_backoff=True,           # 开启指数退避策略
    retry_backoff_max=600,        # 退避上限设为 10 分钟
    retry_jitter=True,            # 增加随机抖动,避免惊群效应
)
def sync_contact_to_crm(self, contact_id: int) -> dict:
    """同步联系人到外部 CRM 系统,网络抖动或临时异常时自动重试。"""
    from apps.crm.services import CRMClient

    try:
        result = CRMClient().sync(contact_id)
        return result
    except CRMClient.RateLimitError as exc:
        # 从响应头提取具体的等待重试秒数
        raise self.retry(exc=exc, countdown=int(exc.retry_after))

幂等任务模式

保证任务在传入相同参数重复执行时依然安全:

@shared_task(name='orders.mark_shipped')
def mark_order_shipped(order_id: int, tracking_number: str) -> None:
    """将订单标记为已发货——支持多次重复执行。"""
    from apps.orders.models import Order

    updated = Order.objects.filter(
        pk=order_id,
        status=Order.Status.PROCESSING,    # 前置防御:仅在处理中状态时更新,避免重复标记
    ).update(
        status=Order.Status.SHIPPED,
        tracking_number=tracking_number,
    )

    if not updated:
        logger.info('mark_order_shipped: order %s already shipped or not found', order_id)

软超时限制任务

from celery.exceptions import SoftTimeLimitExceeded

@shared_task(
    bind=True,
    name='reports.generate_pdf',
    soft_time_limit=120,
    time_limit=150,
)
def generate_pdf_report(self, report_id: int) -> str:
    """生成 PDF 报表,并优雅处理超时情况。"""
    from apps.reports.services import PDFGenerator

    try:
        path = PDFGenerator.build(report_id)
        return path
    except SoftTimeLimitExceeded:
        # 在进程被强行杀死前清理临时文件
        PDFGenerator.cleanup(report_id)
        raise

任务调用方式

from datetime import timedelta
from django.utils import timezone

# 即发即弃(异步派发)
send_welcome_email.delay(user.pk)

# 延迟/定时派发
send_reminder.apply_async(args=[user.pk], countdown=3600)  # 1 小时后执行
send_reminder.apply_async(args=[user.pk], eta=timezone.now() + timedelta(days=1))

# 指定路由队列派发
sync_contact_to_crm.apply_async(args=[contact.pk], queue='high_priority')

# 同步直接执行(仅用于测试或调试环境)
result = generate_pdf_report.apply(args=[report.pk])

Beat 定时任务调度

代码配置调度

# config/settings/base.py
from celery.schedules import crontab

CELERY_BEAT_SCHEDULE = {
    'cleanup-expired-sessions': {
        'task': 'users.cleanup_expired_sessions',
        'schedule': crontab(hour=2, minute=0),   # 每天凌晨 2:00
    },
    'sync-inventory': {
        'task': 'products.sync_inventory',
        'schedule': 60.0,                         # 每 60 秒
    },
    'weekly-digest': {
        'task': 'notifications.send_weekly_digest',
        'schedule': crontab(day_of_week='monday', hour=8, minute=0),
    },
}

数据库动态配置调度(基于 django-celery-beat)

# 可通过 Django Admin 后台或代码动态管理定时任务
from django_celery_beat.models import PeriodicTask, CrontabSchedule
import json

schedule, _ = CrontabSchedule.objects.get_or_create(
    hour='*/6', minute='0',
    timezone='UTC',
)

PeriodicTask.objects.update_or_create(
    name='Sync inventory every 6 hours',
    defaults={
        'crontab': schedule,
        'task': 'products.sync_inventory',
        'args': json.dumps([]),
        'enabled': True,
    }
)

Canvas 工作流:任务串联与并发组合

from celery import chain, group, chord

# Chain(串行):按顺序依次执行,上一步结果自动传给下一个任务
pipeline = chain(
    fetch_data.s(source_id),
    transform_data.s(),          # 接收 fetch_data 的返回值作为首个参数
    load_to_warehouse.s(),
)
pipeline.delay()

# Group(并行):多任务并行并发执行
parallel = group(
    send_welcome_email.s(user_id)
    for user_id in new_user_ids
)
parallel.delay()

# Chord(复合):并行执行任务组,全部完成后触发回调任务
result = chord(
    group(process_chunk.s(chunk) for chunk in data_chunks),
    aggregate_results.s(),       # 回调函数入参为所有分片结果构成的列表
)
result.delay()

错误处理与死信队列(Dead Letter Queue)

# apps/core/tasks.py
from celery.signals import task_failure

@task_failure.connect
def on_task_failure(sender, task_id, exception, args, kwargs, traceback, einfo, **kw):
    """将所有任务失败异常上报至 Sentry 或告警系统。"""
    import sentry_sdk
    with sentry_sdk.new_scope() as scope:
        scope.set_context('celery', {
            'task': sender.name,
            'task_id': task_id,
            'args': args,
            'kwargs': kwargs,
        })
        sentry_sdk.capture_exception(exception)
# 重试达到上限后,将失败任务转存至死信表/队列
@shared_task(
    bind=True,
    max_retries=3,
    name='payments.charge_card',
)
def charge_card(self, order_id: int) -> None:
    from apps.payments.models import Order, FailedCharge

    try:
        _do_charge(order_id)
    except Exception as exc:
        if self.request.retries >= self.max_retries:
            # 持久化至死信表以供人工排查
            FailedCharge.objects.create(
                order_id=order_id,
                error=str(exc),
                task_id=self.request.id,
            )
            return  # 不再抛出异常——标记为永久失败并结束任务
        raise self.retry(exc=exc)

测试 Celery 任务

单元测试(无需消息 Broker)

# tests/test_tasks.py
import pytest
from unittest.mock import patch, MagicMock
from apps.notifications.tasks import send_welcome_email

class TestSendWelcomeEmail:

    @pytest.mark.django_db
    def test_sends_email_to_existing_user(self, user):
        with patch('apps.notifications.services.EmailService') as mock_email:
            send_welcome_email(user.pk)
            mock_email.send_welcome.assert_called_once_with(user)

    @pytest.mark.django_db
    def test_skips_missing_user_gracefully(self):
        """在入队和实际执行间隔期间用户被删除时不应抛错。"""
        send_welcome_email(99999)  # 用户不存在时不抛出异常

集成测试(开启 CELERY_TASK_ALWAYS_EAGER)

# config/settings/test.py
CELERY_TASK_ALWAYS_EAGER = True      # 测试环境中同步直接运行任务
CELERY_TASK_EAGER_PROPAGATES = True  # 任务中的异常直接向上抛出

# tests/test_integration.py
@pytest.mark.django_db
def test_registration_triggers_welcome_email(client):
    with patch('apps.notifications.services.EmailService') as mock_email:
        response = client.post('/api/users/', {
            'email': 'new@example.com',
            'password': 'strongpass123',
        })

    assert response.status_code == 201
    mock_email.send_welcome.assert_called_once()

测试重试逻辑

@pytest.mark.django_db
def test_task_retries_on_connection_error():
    with patch('apps.crm.services.CRMClient.sync') as mock_sync:
        mock_sync.side_effect = ConnectionError('timeout')

        with pytest.raises(ConnectionError):
            sync_contact_to_crm.apply(args=[1], throw=True)

        assert mock_sync.call_count == 1  # 在 eager 同步模式下仅会触发首次尝试

运行监控

# 检查当前活跃的 Worker 和队列状态
celery -A config inspect active
celery -A config inspect stats
celery -A config inspect reserved

# 检查消息队列长度(Redis)
redis-cli llen celery

# Flower:基于 Web 的实时监控面板
pip install flower
celery -A config flower --port=5555

反模式与常见坑点

# 反面教材:直接传递 ORM 模型实例——在任务执行时数据可能已失效
send_welcome_email.delay(user)        # 切勿传递 ORM 模型对象
send_welcome_email.delay(user.pk)     # 始终传递主键 ID

# 反面教材:在生产环境 View 中同步调用任务
result = generate_report.apply()      # 会阻塞当前的请求处理线程

# 反面教材:缺少防重判断的非幂等任务
@shared_task
def charge_and_fulfill(order_id):
    order.charge()     # 任务重试时可能会导致重复扣款!