本文へ移動
cccskills
無料GitHub で公開

celery

Use when running background tasks with Celery - worker and broker configuration (Redis, RabbitMQ), task routing by name vs queue, chains/groups/chords, retry patterns (autoretry_for, retry_backoff), acks_late semantics, failure detection, and monitoring with Flower

インストール方法を見る

含まれるファイル(5)

  • SKILL.md15.4 KB
  • references/beat-routing.md3.8 KB
  • references/calling-workflows.md4.7 KB
  • references/errors-monitoring.md3.3 KB
  • references/testing-performance.md5.6 KB

SKILL.md(原文)

インストールする前に、エージェントに与えられる指示の中身を確認できます。

Celery

Distributed task processing with focus on routing, composition, and failure handling.

When Not to Use Celery

Before adopting Celery, check if you actually need it:

ScenarioBetter Alternative
Short-lived fan-out (5-10 parallel HTTP calls)Simple thread pool or asyncio.gather
Single-server Django app with moderate loadDjango 6.0+ built-in task queue
Work dominated by long IO (10+ min tasks)Task-per-request model or job queue (dramatiq, huey)
Operational overhead unacceptabledjango-tasks, RQ, or simple database polling

Check before adopting:

  • Do you need multiple workers across machines? (If no → simpler queue)
  • Do you have 100+ tasks/minute sustained? (If no → overkill)
  • Can you tolerate operational cost (broker + worker fleet + monitoring)? (If no → huey/dramatiq)

Task Routing and Queue Configuration

Routing by Task Name (Recommended)

# settings.py or celery.py
CELERY_TASK_ROUTES = {
    'myapp.tasks.send_email': {'queue': 'email'},
    'myapp.tasks.process_video': {'queue': 'video'},
    'myapp.tasks.*': {'queue': 'default'},  # Catch-all
}

Why routing by task name beats routing by queue name:

  • Task name is stable; queue name is an implementation detail
  • You can move a task to a different queue without changing call sites
  • Clearer intent: "send_email goes to email queue" vs "this call goes to high_priority"

Precedence: Route vs Queue Argument

# Task decorator sets default queue
@shared_task(queue='default')
def send_email(user_id):
    pass

# apply_async queue overrides decorator
send_email.apply_async(args=[1], queue='high_priority')  # Uses 'high_priority'

# But task_routes overrides both if configured
# CELERY_TASK_ROUTES = {'myapp.tasks.send_email': {'queue': 'critical'}}
# Result: 'critical' wins (routing function is final)

Precedence order: task_routes (routing function) > apply_async(queue=...) > decorator @shared_task(queue=...)

Per-Queue Worker Startup

# High-priority queue: solo pool for immediate processing
celery -A myproject worker -Q high -P solo -n high@%h

# Low-priority queue: prefork for throughput
celery -A myproject worker -Q low -P prefork --concurrency=4 -n low@%h

# Mixed queues: one worker handles both
celery -A myproject worker -Q high,low -P prefork --concurrency=2

Pool choices:

  • solo: Single process, no concurrency. Use for debugging or tasks that can't be parallelized.
  • prefork: Default. Good for CPU-bound or mixed workloads.
  • eventlet/gevent: High concurrency for I/O-bound tasks (100+ simultaneous).

Critical Failure: Route to Queue with No Worker

# Task routed to queue no one consumes
CELERY_TASK_ROUTES = {'myapp.tasks.cleanup': {'queue': 'cleanup'}}

# Problem: No worker started with -Q cleanup
# Result: Task sits in broker forever, silently. No error, no retry.

Detection:

# Check which queues have active workers
celery -A myproject inspect active

# Check queue depth
redis-cli LLEN celery  # For Redis broker

# Flower UI: Monitor queue lengths per queue

Fix: Start a worker for the queue OR remove the route so tasks go to a consumed queue.

Task Composition: Chains, Groups, Chords

When to Use Each

PatternUse CaseRequires Backend
chainSequential steps (A → B → C)No (but needed for .get())
groupParallel independent tasksNo (but needed for results)
chordParallel + callback (wait for all, then run)Yes (stores intermediate results)

Chains (Sequential)

from celery import chain

# Correct: Results flow through
workflow = chain(
    validate_user.s(user_id),
    process_data.s(),
    send_notification.s()
)
result = workflow.apply_async()
# Each task's output becomes next task's input

# Error handling: link error callback
from celery import group

error_handler = handle_error.s()
workflow = (
    chain(validate.s(), process.s()) |
    save_result.s()
)
workflow.apply_async(link_error=error_handler)

Groups (Parallel)

from celery import group

# Correct: Independent parallel tasks
job = group(
    fetch_data.s(source_id)
    for source_id in [1, 2, 3, 4, 5]
)
result = job.apply_async()
all_results = result.get()  # [res1, res2, res3, res4, res5]

# Wrong: Passing unbounded group into chord (deadlock risk)
# See chord section below

Chords (Parallel + Callback)

from celery import chord, group

# Correct form: Bounded group as chord header
result = chord(
    group(
        process_item.s(item_id)
        for item_id in item_ids[:100]  # Bounded list
    ),
    summarize_results.s()  # Callback receives [res1, res2, ...]
).apply_async()

# Incorrect form that hangs: Unbounded generator
# DO NOT DO THIS:
# result = chord(
#     group(process_item.s(item_id) for item_id in infinite_generator()),
#     callback.s()
# )
# Why: Chord waits for ALL header tasks to complete. Generator never ends.

# Also incorrect: Group with dynamic size > worker capacity
# If you have 4 workers but group(1000 tasks), some tasks wait.
# Chord callback never fires until all 1000 complete.
# Solution: Use chunks or multiple chords

Chord backend requirement:

# settings.py
CELERY_RESULT_BACKEND = 'redis://localhost:6379/0'  # Required for chords

# Without backend, chord callback never fires
# Error: Chord 'xxx' raised: 'NoneType' object is not subscriptable

Retry and Failure Handling

Auto-Retry Configuration

@shared_task(
    bind=True,
    autoretry_for=(ConnectionError, TimeoutError, ExternalAPIError),
    retry_backoff=True,          # Exponential backoff (1s, 2s, 4s, 8s...)
    retry_backoff_max=600,       # Cap at 10 minutes
    retry_jitter=True,           # Add randomness to prevent thundering herd
    max_retries=3,
)
def fetch_external_data(self, url):
    """Auto-retry on transient failures."""
    response = requests.get(url, timeout=10)
    response.raise_for_status()
    return response.json()

Retry parameters explained:

  • autoretry_for: Tuple of exception classes that trigger auto-retry
  • retry_backoff=True: Delays grow exponentially (default: 1s to 60s)
  • retry_backoff_max: Maximum delay between retries
  • retry_jitter=True: Add ±25% randomness to each delay
  • max_retries: Total retry attempts (excluding first execution)

Manual Retry

@shared_task(bind=True, max_retries=5)
def process_with_manual_retry(self, record_id):
    try:
        record = Record.objects.get(id=record_id)
        return record.process()
    except Record.DoesNotExist:
        # Don't retry: data missing permanently
        logger.error(f"Record {record_id} not found")
        return None
    except DatabaseError as exc:
        # Retry with custom countdown
        countdown = min(60 * (2 ** self.request.retries), 300)
        raise self.retry(exc=exc, countdown=countdown)
    except Exception as exc:
        # Unexpected error: re-raise without retry
        logger.exception("Unexpected error")
        raise

What NOT to Retry

# Non-idempotent side effects: Don't retry
@shared_task(bind=True, max_retries=3)
def charge_credit_card(self, payment_id):
    """WRONG: Charging card multiple times on retry."""
    payment = Payment.objects.get(id=payment_id)
    stripe.Charge.create(amount=payment.amount, card=payment.card_id)
    # If worker crashes after charge but before DB update, retry charges again

# Correct: Make idempotent
@shared_task(bind=True, max_retries=3)
def charge_credit_card_idempotent(self, payment_id):
    payment = Payment.objects.select_for_update().get(id=payment_id)
    
    if payment.status == 'charged':
        return {'status': 'already_charged'}
    
    with transaction.atomic():
        payment.refresh_from_db()
        if payment.status == 'charged':
            return {'status': 'already_charged'}
        
        stripe.Charge.create(amount=payment.amount, card=payment.card_id)
        payment.status = 'charged'
        payment.save()
    
    return {'status': 'charged'}

acks_late and task_reject_on_worker_lost

# settings.py
CELERY_ACKS_LATE = True  # Acknowledge task AFTER execution completes

CELERY_TASK_REJECT_ON_WORKER_LOST = True  # Reject unacked tasks when worker dies

Semantics and trade-offs:

SettingBehaviorTrade-off
acks_late=False (default)Ack immediately on receiptWorker crash = task lost (no retry)
acks_late=TrueAck after task completesWorker crash = task requeued (may run twice)
reject_on_worker_lost=TrueReject unacked tasks on worker deathSame as acks_late but explicit

At-least-once delivery means idempotency is your responsibility:

# With acks_late=True, this task may run twice:
@shared_task(acks_late=True)
def send_welcome_email(user_id):
    # WRONG: Sends duplicate email
    User.objects.get(id=user_id).send_email()
    
# Correct: Idempotent
@shared_task(acks_late=True)
def send_welcome_email_idempotent(user_id):
    user = User.objects.get(id=user_id)
    if user.welcome_email_sent:
        return {'status': 'already_sent'}
    
    with transaction.atomic():
        user.refresh_from_db()
        if user.welcome_email_sent:
            return {'status': 'already_sent'}
        user.send_email()
        user.welcome_email_sent = True
        user.save()

Common Issues (Symptom-First)

Symptom: Tasks Never Execute (Stuck in PENDING)

How to confirm:

# 1. Check worker is running and consuming correct queue
celery -A myproject inspect active

# 2. Check queue routing
python -c "from myproject.celery import app; print(app.conf.task_routes)"

# 3. Check broker queue depth
redis-cli LLEN celery  # Should decrease as tasks execute

# 4. Verify task name matches
celery -A myproject inspect registered  # Task must appear here

Fix:

  • Queue misrouting: Add route or start worker with -Q <queue_name>
  • Missing default_queue: Set CELERY_TASK_DEFAULT_QUEUE = 'default'
  • Task not registered: Check autodiscover_tasks() is called

Symptom: Tasks Silently Lost (No Error, No Execution)

How to confirm:

# 1. Check broker/result backend mismatch
# If broker=redis://localhost:6379/0 but result_backend=redis://localhost:6379/1
# Results go to different Redis instance

# 2. Check worker logs for "Task xxx raised" without traceback
# Worker may have crashed before ack

# 3. Verify task version probe
celery -A myproject inspect registered  # Compare task versions across workers

Fix:

  • Ensure broker_url and result_backend point to same broker instance
  • Set CELERY_TASK_ALWAYS_EAGER = False (debug mode default)
  • Check firewall/network between workers and broker

Symptom: Stale Workers Running Old Code

How to confirm:

# In task, probe version
@shared_task
def version_probe():
    import myapp.tasks
    return myapp.tasks.__version__  # Or check git commit hash

# Compare across workers
celery -A myproject inspect active

Fix:

  • Restart workers after code deploy: systemctl restart celery
  • Use --max-tasks-per-child=1000 to force periodic restarts
  • Deploy with zero-downtime: start new workers, drain old workers

Symptom: Worker Killed Mid-Task (Data Corruption)

How to confirm:

# Check for OOM killer in system logs
dmesg | grep -i "out of memory"
grep "Out of memory" /var/log/syslog

# Check worker process exit codes
journalctl -u celery | grep "killed"

Fix:

  • Set acks_late=True + idempotent tasks (task re-runs on crash)
  • Increase worker memory limits
  • Use CELERY_WORKER_MAX_MEMORY_PER_CHILD = 400000 (KB)

Symptom: Time-Limit Exceptions

# Hard time limit: Worker kills task immediately
@shared_task(time_limit=300)  # 5 minutes hard limit
def long_task():
    pass

# Soft time limit: Task receives SoftTimeLimitExceeded exception
@shared_task(soft_time_limit=240)  # 4 minutes soft limit
def graceful_task():
    try:
        return do_work()
    except SoftTimeLimitExceeded:
        # Cleanup before forced termination
        cleanup()
        raise

Detection:

# Check worker logs for "Task raised SoftTimeLimitExceeded"
grep "SoftTimeLimitExceeded" /var/log/celery/worker.log

Symptom: Results Never Appear

How to confirm:

# 1. Check result backend is configured
print(app.conf.result_backend)  # Must not be None

# 2. Check backend matches between sender and receiver
# If task sent with backend=redis://... but result.get() uses different backend

# 3. Check result expiration
print(app.conf.result_expires)  # Default 1 day

Fix:

  • Set CELERY_RESULT_BACKEND = 'redis://localhost:6379/0'
  • Don't set ignore_result=True unless you don't need results
  • Increase result_expires for long-running workflows

Symptom: Duplicate Task Execution

How to confirm:

# Check acks_late setting
celery -A myproject inspect config | grep ACKS_LATE

# Check for worker restarts during task execution
journalctl -u celery | grep -E "(restarting|killed|OOM)"

Fix:

  • Make tasks idempotent (see retry section)
  • Use database unique constraints
  • Implement deduplication tokens

Best Practices

Idempotent Tasks

@shared_task(bind=True)
def process_payment(self, payment_id):
    """Idempotent: Safe to run multiple times."""
    payment = Payment.objects.select_for_update().get(id=payment_id)
    
    if payment.status == 'completed':
        return {'status': 'already_processed'}
    
    with transaction.atomic():
        payment.refresh_from_db()
        if payment.status == 'completed':
            return {'status': 'already_processed'}
        
        result = payment.charge()
        payment.status = 'completed'
        payment.save()
    
    return result

Task Granularity

# Bad: Monolithic task
@shared_task
def process_order_bad(order_id):
    order = Order.objects.get(id=order_id)
    order.validate()
    order.charge()
    order.ship()

# Good: Composable tasks
@shared_task
def validate_order(order_id):
    Order.objects.get(id=order_id).validate()
    charge_order.delay(order_id)

@shared_task
def charge_order(order_id):
    order = Order.objects.get(id=order_id)
    order.charge()
    ship_order.delay(order_id)

Deep Dives

Load these on demand for detailed reference:

  • Calling Workflows (references/calling-workflows.md) — apply_async options, signatures, chains, groups, chords, chunks, and task states
  • Periodic Tasks & Routing (references/beat-routing.md) — Celery Beat scheduling, crontab syntax, django-celery-beat, queue routing
  • Error Handling & Monitoring (references/errors-monitoring.md) — Retries, error callbacks, DLQ, Flower, CLI monitoring, Prometheus
  • Testing & Performance (references/testing-performance.md) — pytest testing, worker configuration, task optimization, memory management

References

レビュー

まだレビューはありません。使ってみた感想をお寄せください。

同じリポジトリのスキル

概要と使いどころ

aiohttp

無料

Use when building Python async HTTP services or clients with aiohttp - web server routing, middleware, WebSocket, SSE, streaming, client sessions, pytest-aiohttp testing, or troubleshooting SSL and timeout issues

日本語の概要は準備中です。原文の説明を表示しています。

CodeAtCode/oss-ai-skills222026年10月9日 更新

ast-grep

無料

Use when doing structural code search and rewriting - ast-grep linting, refactoring, multi-language patterns

日本語の概要は準備中です。原文の説明を表示しています。

CodeAtCode/oss-ai-skills222026年10月9日 更新

Use when building GBA games with the BPCore Lua engine - entity, sprite and tilemap functions, SRAM save and load, link cable multiplayer protocol, camera and scrolling, or optimization patterns

日本語の概要は準備中です。原文の説明を表示しています。

CodeAtCode/oss-ai-skills222026年10月9日 更新

django

無料

Use when building Django applications - security hardening, authentication and permissions, ORM optimization, PostgreSQL features, Django 6.0, migrations, testing, and ecosystem libraries

日本語の概要は準備中です。原文の説明を表示しています。

CodeAtCode/oss-ai-skills222026年10月9日 更新

Use when customizing Django Admin - save_formset, get_search_results, formsets, queryset optimization, db_index, custom URLs

日本語の概要は準備中です。原文の説明を表示しています。

CodeAtCode/oss-ai-skills222026年10月9日 更新

Use when implementing Django authentication - local accounts, OAuth, email verification, MFA, OIDC, django-organizations

日本語の概要は準備中です。原文の説明を表示しています。

CodeAtCode/oss-ai-skills222026年10月9日 更新

CodeAtCode のスキルをすべて見る

このスキルの問題を報告する