SKILL.md
唯讀
名稱
django-celery
描述
Django + Celery 非同步任務設計模式 — 涵蓋環境設定、任務設計、Beat 排程、重試機制、Canvas 工作流、監控與測試。適用於為 Django 應用程式新增背景作業、定時任務或非同步處理時。
Django + Celery 非同步任務設計模式
在 Django 中使用 Celery 搭配 Redis 或 RabbitMQ 進行背景任務處理的生產級最佳實踐。
何時啟用
- 為 Django 應用程式新增背景作業或非同步處理
- 實作週期性 / 定時任務
- 將耗時操作(發送 Email、生成 PDF、呼叫外部 API)移出 HTTP 請求週期
- 設定 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 # 軟性時間限制:拋出 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('Welcome email sent to user %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 點
},
'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),
},
}
資料庫動態排程 (via django-celery-beat)
# 透過 Django 後台或程式碼動態管理週期性任務
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()
錯誤處理與死信佇列
# 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
反模式
# 錯誤:傳遞 Model 實例 — 任務執行時資料可能已經過時
send_welcome_email.delay(user) # 切勿直接傳遞 ORM 物件
send_welcome_email.delay(user.pk) # 一律傳入主鍵 (PK)
# 錯誤:在生產環境的 View 中同步呼叫任務
result = generate_report.apply() # 這會阻塞整個請求執行緒
# 錯誤:缺乏防護機制的非冪等任務
@shared_task
def charge_and_fulfill(order_id):
order.charge() # 若任務發生重試,可能會重複扣款!




