Skip to main content Skills Marketplace Découvrez et explorez les compétences IA créées par la communauté.
Installer avec Codex ou Claude Copiez ce prompt, collez-le dans Codex, Claude ou un autre assistant, puis laissez-le vérifier la page du skill et l'installer pour vous.
Copier le promptAfficher les détails du prompt Une commande directe contourne le prompt de vérification. Examinez la source avant de l'exécuter.
npx skills add https://github.com/tomevault-io/skills-registry --skill django-celeryLa commande reste sur une seule ligne. Faites défiler horizontalement pour la vérifier avant de la copier.
Vous préférez une copie locale ? Téléchargez les fichiers actuellement disponibles dans SkillsMP.
Télécharger Zip Téléchargement... Explorateur de fichiers
2 fichiers Métiers associés SOC
Basé sur la classification professionnelle SOC
name django-celery description Django + Celery 异步任务模式 — 配置、任务设计、Beat 调度、重试、Canvas 工作流、监控和测试。当向 Django 应用添加后台任务、定时任务或异步处理时使用。 Use when this capability is needed. metadata {"author":"aaione"}
Django + Celery 异步任务模式
在 Django 中使用 Celery 配合 Redis 或 RabbitMQ 进行后台任务处理的生产级模式。
何时激活
向 Django 应用添加后台任务或异步处理
实现周期性/定时任务
从请求周期中卸载慢操作(邮件、PDF 生成、API 调用)
设置 Celery Beat 进行类 cron 调度
调试任务失败、重试或队列积压
为 Celery 任务编写测试
项目设置
安装
pip install celery[redis] django-celery-results django-celery-beat
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()
@app.task(bind=True , ignore_result=True )
def debug_task (self ):
( )
print
f'请求:{self.request!r} '
from .celery import app as celery_app
__all__ = ('celery_app' ,)
Django 设置
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
CELERY_TASK_SOFT_TIME_LIMIT = 25 * 60
CELERY_WORKER_PREFETCH_MULTIPLIER = 1
CELERY_TASK_ACKS_LATE = True
CELERY_RESULT_EXPIRES = 60 * 60 * 24
CELERY_BEAT_SCHEDULER = 'django_celery_beat.schedulers:DatabaseScheduler'
INSTALLED_APPS += [
'django_celery_results' ,
'django_celery_beat' ,
]
运行 Worker
celery -A config worker --loglevel=info
celery -A config beat --loglevel=info --scheduler django_celery_beat.schedulers:DatabaseScheduler
celery -A config worker --beat --loglevel=info
celery -A config worker --loglevel=warning --concurrency=4 -Q default,high_priority
任务设计模式
基本任务
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: 用户 %s 未找到' , 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 ,
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: 订单 %s 已发货或未找到' , 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 )
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 调度(周期性任务)
代码定义的调度
from celery.schedules import crontab
CELERY_BEAT_SCHEDULE = {
'cleanup-expired-sessions' : {
'task' : 'users.cleanup_expired_sessions' ,
'schedule' : crontab(hour=2 , minute=0 ),
},
'sync-inventory' : {
'task' : 'products.sync_inventory' ,
'schedule' : 60.0 ,
},
'weekly-digest' : {
'task' : 'notifications.send_weekly_digest' ,
'schedule' : crontab(day_of_week='monday' , hour=8 , minute=0 ),
},
}
数据库定义的调度(通过 django-celery-beat)
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='每 6 小时同步库存' ,
defaults={
'crontab' : schedule,
'task' : 'products.sync_inventory' ,
'args' : json.dumps([]),
'enabled' : True ,
}
)
Canvas:链式和分组任务 from celery import chain, group, chord
pipeline = chain(
fetch_data.s(source_id),
transform_data.s(),
load_to_warehouse.s(),
)
pipeline.delay()
parallel = group(
send_welcome_email.s(user_id)
for user_id in new_user_ids
)
parallel.delay()
result = chord(
group(process_chunk.s(chunk) for chunk in data_chunks),
aggregate_results.s(),
)
result.delay()
错误处理和死信队列
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)
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 进行集成测试
CELERY_TASK_ALWAYS_EAGER = True
CELERY_TASK_EAGER_PROPAGATES = True
@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
监控
celery -A config inspect active
celery -A config inspect stats
celery -A config inspect reserved
redis-cli llen celery
pip install flower
celery -A config flower --port=5555
反模式
send_welcome_email.delay(user)
send_welcome_email.delay(user.pk)
result = generate_report.apply()
@shared_task
def charge_and_fulfill (order_id ):
order.charge()
order.fulfill()
@shared_task
def charge_and_fulfill (order_id ):
order = Order.objects.select_for_update().get(pk=order_id)
if order.status != Order.Status.PENDING:
return
order.charge()
order.fulfill()
生产检查清单 检查项 设置 Worker 崩溃时自动重启 supervisord 或 systemd unitCELERY_TASK_ACKS_LATE = TrueWorker 崩溃时重新入队任务 CELERY_WORKER_PREFETCH_MULTIPLIER = 1长任务的公平分配 按优先级分离队列 -Q default,high_priority,low_priority设置 CELERY_TASK_SOFT_TIME_LIMIT 硬终止前的优雅超时 Sentry 集成 捕获所有 task_failure 信号 Flower 或其他监控工具 队列深度可视化 Beat 仅在单节点运行 防止重复执行定时任务
相关技能
django-patterns — ORM、服务层和项目结构
django-tdd — 测试 Django 模型、视图和服务
python-testing — pytest 配置和 fixture