[OK] CHEATSHEET CELERY PYTHON - GUIDE COMPLET
[OK] Celery : système de file d'attente de tâches asynchrones distribuées
# Permet d'exécuter des tâches en arrière-plan (envoi emails, traitement images, etc.)

[OK] 1. INSTALLATION


# Installation de base
# pip install celery

# Avec Redis comme broker
# pip install celery[redis]

# Avec RabbitMQ comme broker
# pip install celery[amqp]

# Avec base de données comme backend
# pip install celery[sqlalchemy]

# Installation complète
# pip install "celery[redis,sqlalchemy,msgpack]"

[OK] 2. CONFIGURATION DE BASE


# Fichier: celery_app.py
from celery import Celery

# Créer instance Celery avec Redis comme broker
app = Celery('myapp',
             broker='redis://localhost:6379/0',      # Message broker
             backend='redis://localhost:6379/0')      # Results backend

# Avec RabbitMQ
app = Celery('myapp',
             broker='amqp://guest@localhost//',
             backend='rpc://')

# Configuration complète
app = Celery('myapp')
app.conf.update(
    broker_url='redis://localhost:6379/0',
    result_backend='redis://localhost:6379/0',
    task_serializer='json',                          # Format de sérialisation
    accept_content=['json'],                         # Formats acceptés
    result_serializer='json',
    timezone='Europe/Paris',                         # Timezone
    enable_utc=True,
)

[OK] 3. CRÉER DES TÂCHES


# Tâche simple avec décorateur
@app.task
def add(x, y):
    """Additionne deux nombres"""
    return x + y

# Tâche avec nom personnalisé
@app.task(name='tasks.multiply')
def multiply(x, y):
    """Multiplie deux nombres"""
    return x * y

# Tâche avec options
@app.task(
    name='tasks.send_email',
    bind=True,                    # Accès à l'instance de tâche (self)
    max_retries=3,                # Nombre maximum de tentatives
    soft_time_limit=60,           # Timeout souple (exception)
    time_limit=120,               # Timeout dur (kill process)
    ignore_result=False,          # Stocker le résultat
    track_started=True,           # Tracer le démarrage
    acks_late=True,               # Acquitter après exécution
)
def send_email(self, to, subject, body):
    """Envoie un email"""
    try:
        # Code d'envoi d'email
        return f"Email envoyé à {to}"
    except Exception as exc:
        # Réessayer après 5 secondes
        raise self.retry(exc=exc, countdown=5)

# Tâche avec classe
from celery import Task

class DatabaseTask(Task):
    """Tâche avec connexion DB"""
    _db = None
    
    @property
    def db(self):
        if self._db is None:
            # Initialiser connexion DB
            self._db = connect_to_database()
        return self._db

@app.task(base=DatabaseTask)
def process_data(data):
    """Traite des données avec DB"""
    db = process_data.db
    # Utiliser la connexion
    return "Données traitées"

[OK] 4. APPELER DES TÂCHES


# Exécution asynchrone (immédiate)
result = add.delay(4, 6)
# Équivalent à:
result = add.apply_async(args=[4, 6])

# Exécution avec options
result = add.apply_async(
    args=[4, 6],
    kwargs={'x': 4, 'y': 6},
    countdown=10,              # Exécuter dans 10 secondes
    eta=datetime.utcnow() + timedelta(hours=1),  # À une heure précise
    expires=300,               # Expire dans 5 minutes
    retry=True,                # Autoriser retry
    retry_policy={
        'max_retries': 3,
        'interval_start': 0,
        'interval_step': 0.2,
        'interval_max': 0.2,
    },
    queue='high_priority',     # File spécifique
    routing_key='tasks.add',
    priority=9,                # Priorité (0-9, plus haut = plus prioritaire)
)

# Exécution différée (ETA)
from datetime import datetime, timedelta
eta = datetime.utcnow() + timedelta(hours=2)
result = add.apply_async(args=[4, 6], eta=eta)

# Exécution dans X secondes
result = add.apply_async(args=[4, 6], countdown=60)

# Lien de tâches (chaînage)
# Exécuter add, puis multiply avec le résultat
from celery import chain
result = chain(add.s(4, 6), multiply.s(2))()
# 4 + 6 = 10, puis 10 * 2 = 20

# Groupe de tâches (parallèle)
from celery import group
job = group([
    add.s(2, 2),
    add.s(4, 4),
    add.s(8, 8)
])
result = job.apply_async()

# Chord (groupe + callback)
from celery import chord
callback = multiply.s(2)
result = chord([add.s(2, 2), add.s(4, 4)])(callback)
# (2+2) + (4+4) = 12, puis 12 * 2 = 24

# Map (appliquer à une liste)
from celery import group
result = group(add.s(i, i) for i in range(10))()

# Chunks (diviser en morceaux)
result = add.chunks(zip(range(100), range(100)), 10)()
# Divise 100 tâches en chunks de 10

[OK] 5. RÉSULTATS DE TÂCHES


# Obtenir le résultat (bloquant)
result = add.delay(4, 6)
value = result.get(timeout=10)  # Attend max 10 secondes
print(value)  # 10

# Vérifier l'état
result.ready()      # True si terminé
result.successful() # True si succès
result.failed()     # True si échec
result.state        # 'PENDING', 'STARTED', 'SUCCESS', 'FAILURE', 'RETRY'

# Informations sur le résultat
result.id           # ID unique de la tâche
result.task_id      # Même chose
result.status       # État actuel
result.result       # Résultat ou exception
result.traceback    # Traceback si erreur

# Attendre sans bloquer
import time
while not result.ready():
    print("En cours...")
    time.sleep(1)
print(result.get())

# Obtenir résultat avec gestion d'erreur
try:
    value = result.get(timeout=10, propagate=True)
except Exception as exc:
    print(f"Erreur: {exc}")

# Révoquer une tâche
result.revoke(terminate=True)  # Arrêter immédiatement

# Oublier le résultat (libérer mémoire)
result.forget()

[OK] 6. GESTION D'ERREURS ET RETRY


# Retry automatique
@app.task(bind=True, max_retries=3)
def unreliable_task(self, data):
    try:
        # Code qui peut échouer
        process_data(data)
    except Exception as exc:
        # Réessayer avec délai exponentiel
        raise self.retry(exc=exc, countdown=2 ** self.request.retries)

# Retry avec conditions
@app.task(bind=True, autoretry_for=(ConnectionError,), retry_kwargs={'max_retries': 5})
def fetch_data(self, url):
    response = requests.get(url)
    return response.json()

# Gérer les échecs
@app.task(bind=True, on_failure=log_failure)
def risky_task(self):
    # Code risqué
    pass

def log_failure(self, exc, task_id, args, kwargs, einfo):
    print(f"Tâche {task_id} échouée: {exc}")

[OK] 7. TÂCHES PÉRIODIQUES (BEAT)


from celery.schedules import crontab

# Configuration dans app
app.conf.beat_schedule = {
    # Exécuter toutes les 30 secondes
    'add-every-30-seconds': {
        'task': 'tasks.add',
        'schedule': 30.0,
        'args': (16, 16)
    },
    
    # Exécuter tous les lundis à 7h30
    'send-report-monday': {
        'task': 'tasks.send_weekly_report',
        'schedule': crontab(hour=7, minute=30, day_of_week=1),
    },
    
    # Tous les jours à minuit
    'cleanup-daily': {
        'task': 'tasks.cleanup',
        'schedule': crontab(hour=0, minute=0),
    },
    
    # Toutes les heures
    'sync-hourly': {
        'task': 'tasks.sync_data',
        'schedule': crontab(minute=0),
    },
    
    # Tous les premiers du mois
    'monthly-report': {
        'task': 'tasks.monthly_report',
        'schedule': crontab(day_of_month=1, hour=9, minute=0),
    },
}

# Exemples de crontab
crontab(minute=0, hour='*/3')           # Toutes les 3 heures
crontab(minute=0, hour='0,3,6,9,12,15,18,21')  # Heures spécifiques
crontab(day_of_week='mon,wed,fri')      # Certains jours
crontab(minute='*/15')                  # Toutes les 15 minutes
crontab(hour=7, minute=30, day_of_week='1-5')  # Lun-Ven à 7h30

# Démarrer beat scheduler
# celery -A celery_app beat --loglevel=info

[OK] 8. WORKERS


# Démarrer un worker
# celery -A celery_app worker --loglevel=info

# Worker avec options
# celery -A celery_app worker \
#     --loglevel=info \
#     --concurrency=4 \          # Nombre de processus
#     --pool=prefork \           # Type de pool (prefork, gevent, eventlet)
#     --queue=high_priority \    # Queue spécifique
#     --hostname=worker1@%h \    # Nom du worker
#     --max-tasks-per-child=100  # Redémarrer après N tâches

# Pool types
# prefork  : Processus multiples (CPU-bound)
# gevent   : Coroutines (I/O-bound)
# eventlet : Coroutines (I/O-bound)
# solo     : Single thread (debug)

# Worker avec autoscale
# celery -A celery_app worker --autoscale=10,3
# Min 3 workers, max 10 selon charge

[OK] 9. QUEUES ET ROUTING


# Configuration des queues
app.conf.task_routes = {
    'tasks.send_email': {'queue': 'emails'},
    'tasks.process_image': {'queue': 'media'},
    'tasks.*': {'queue': 'default'},
}

# Ou avec routing explicite
from kombu import Exchange, Queue

app.conf.task_queues = (
    Queue('default', Exchange('default'), routing_key='default'),
    Queue('high_priority', Exchange('high'), routing_key='high'),
    Queue('low_priority', Exchange('low'), routing_key='low'),
)

# Router personnalisé
class MyRouter:
    def route_for_task(self, task, args=None, kwargs=None):
        if task.startswith('tasks.high'):
            return {'queue': 'high_priority'}
        return {'queue': 'default'}

app.conf.task_routes = (MyRouter(),)

# Envoyer à une queue spécifique
result = add.apply_async(args=[4, 6], queue='high_priority')

# Worker écoutant plusieurs queues
# celery -A celery_app worker -Q high_priority,default

[OK] 10. MONITORING ET INSPECTION


# API Inspect
from celery import current_app

inspect = current_app.control.inspect()

# Workers actifs
workers = inspect.active()
print(workers)  # {'worker1@host': [...], ...}

# Tâches actives
active_tasks = inspect.active()

# Tâches planifiées
scheduled = inspect.scheduled()

# Tâches réservées
reserved = inspect.reserved()

# Statistiques
stats = inspect.stats()

# Ping workers
ping = inspect.ping()

# Registered tasks
registered = inspect.registered()

# Contrôle des workers
from celery import current_app

control = current_app.control

# Révoquer une tâche
control.revoke(task_id, terminate=True)

# Arrêter un worker
control.shutdown()

# Ajouter consumer à une queue
control.add_consumer('high_priority')

# Annuler consumer
control.cancel_consumer('low_priority')

# Rate limit
control.rate_limit('tasks.process_image', '10/m')  # 10 par minute

[OK] 11. ÉVÉNEMENTS ET MONITORING


# Activer événements
# celery -A celery_app worker --loglevel=info -E

# Écouter les événements
from celery import Celery

def on_task_sent(event):
    print(f"Tâche envoyée: {event['uuid']}")

def on_task_success(event):
    print(f"Tâche réussie: {event['uuid']}, résultat: {event['result']}")

def on_task_failure(event):
    print(f"Tâche échouée: {event['uuid']}, exception: {event['exception']}")

# Monitorer les événements
with app.connection() as connection:
    receiver = app.events.Receiver(connection, handlers={
        'task-sent': on_task_sent,
        'task-succeeded': on_task_success,
        'task-failed': on_task_failure,
    })
    receiver.capture(limit=None, timeout=None, wakeup=True)

# Flower (Web UI pour monitoring)
# pip install flower
# celery -A celery_app flower --port=5555
# Accéder via http://localhost:5555

[OK] 12. CONFIGURATION AVANCÉE


# Configuration complète
app.conf.update(
    # Broker
    broker_url='redis://localhost:6379/0',
    broker_connection_retry_on_startup=True,
    broker_connection_retry=True,
    broker_connection_max_retries=10,
    
    # Results
    result_backend='redis://localhost:6379/0',
    result_expires=3600,              # Résultats expirent après 1h
    result_persistent=True,           # Persister les résultats
    result_extended=True,             # Informations étendues
    
    # Serialization
    task_serializer='json',
    result_serializer='json',
    accept_content=['json'],
    
    # Timezone
    timezone='Europe/Paris',
    enable_utc=True,
    
    # Task execution
    task_always_eager=False,          # True = exécution synchrone (debug)
    task_eager_propagates=True,
    task_ignore_result=False,
    task_track_started=True,
    task_acks_late=True,              # Acquitter après exécution
    task_reject_on_worker_lost=True,
    
    # Performance
    worker_prefetch_multiplier=4,     # Nombre de tâches préchargées
    worker_max_tasks_per_child=1000,  # Redémarrer worker après N tâches
    worker_disable_rate_limits=False,
    
    # Time limits
    task_soft_time_limit=60,          # Timeout souple
    task_time_limit=120,              # Timeout dur
    
    # Logging
    worker_hijack_root_logger=False,
    worker_log_format='[%(asctime)s: %(levelname)s/%(processName)s] %(message)s',
)

[OK] 13. SIGNATURES (PRIMITIVES)


# Signature simple
sig = add.signature(args=(4, 6), countdown=10)
# Ou raccourci:
sig = add.s(4, 6)

# Partial signature (arguments incomplets)
partial = multiply.s(2)  # y manquant
result = partial.apply_async(args=[10])  # 10 * 2 = 20

# Chaîne (chain)
from celery import chain
workflow = chain(
    add.s(4, 4),
    multiply.s(2),
    add.s(10)
)
result = workflow()  # (4+4)*2+10 = 26

# Groupe (group) - parallèle
from celery import group
job = group(add.s(i, i) for i in range(10))
result = job()

# Chord - groupe + callback
from celery import chord
workflow = chord([
    add.s(2, 2),
    add.s(4, 4),
    add.s(8, 8)
])(multiply.s(2))
# (2+2) + (4+4) + (8+8) = 4 + 8 + 16 = 28, puis 28*2 = 56

# Map - appliquer à une liste
from celery import group
job = add.map([(i, i) for i in range(10)])
result = job.apply_async()

# Starmap - unpack arguments
from celery import group
job = add.starmap(zip(range(10), range(10)))

# Chunks - diviser en morceaux
job = add.chunks(zip(range(100), range(100)), 10)
# 100 tâches divisées en 10 chunks de 10

[OK] 14. CANVAS - WORKFLOWS COMPLEXES


# Workflow complexe
from celery import chain, group, chord

# Étape 1: Traiter en parallèle
step1 = group(
    process_data.s(data1),
    process_data.s(data2),
    process_data.s(data3)
)

# Étape 2: Agréger résultats
step2 = aggregate_results.s()

# Étape 3: Envoyer rapport
step3 = send_report.s()

# Composer le workflow
workflow = chain(
    chord(step1, step2),
    step3
)

# Exécuter
result = workflow.apply_async()

# Workflow conditionnel (avec immutability)
@app.task
def conditional_task(value):
    if value > 10:
        return process_large.si(value).apply_async()
    else:
        return process_small.si(value).apply_async()

[OK] 15. TESTS


# Mode eager (synchrone) pour tests
app.conf.task_always_eager = True
app.conf.task_eager_propagates = True

# Test unitaire
def test_add():
    result = add.delay(4, 6)
    assert result.get() == 10

# Mock avec pytest
import pytest
from unittest.mock import patch

def test_send_email():
    with patch('tasks.send_email.delay') as mock_task:
        send_email.delay('user@example.com', 'Test', 'Body')
        mock_task.assert_called_once()

# Test avec backend de test
from celery import current_app
from celery.contrib.testing.worker import start_worker

@pytest.fixture(scope='session')
def celery_worker():
    with start_worker(current_app, perform_ping_check=False):
        yield

[OK] 16. EXEMPLES PRATIQUES


# 1. Envoi d'emails en masse
@app.task
def send_bulk_emails(email_list):
    for email in email_list:
        send_email.delay(email['to'], email['subject'], email['body'])

# 2. Traitement d'images
@app.task(bind=True, max_retries=3)
def process_image(self, image_path):
    try:
        img = Image.open(image_path)
        # Redimensionner, optimiser, etc.
        img.thumbnail((800, 600))
        img.save(image_path, optimize=True)
        return image_path
    except Exception as exc:
        raise self.retry(exc=exc, countdown=60)

# 3. Export de données
@app.task
def export_to_csv(query_params):
    data = fetch_data_from_db(query_params)
    filename = f"export_{datetime.now():%Y%m%d_%H%M%S}.csv"
    with open(filename, 'w') as f:
        writer = csv.writer(f)
        writer.writerows(data)
    return filename

# 4. Nettoyage périodique
@app.task
def cleanup_old_files():
    cutoff = datetime.now() - timedelta(days=30)
    for file in os.listdir('/tmp/uploads'):
        filepath = os.path.join('/tmp/uploads', file)
        if os.path.getmtime(filepath) < cutoff.timestamp():
            os.remove(filepath)

# 5. Web scraping
@app.task(rate_limit='10/m')  # Max 10 par minute
def scrape_url(url):
    response = requests.get(url)
    # Parser et extraire données
    return parse_html(response.text)

# 6. Génération de rapports
@app.task
def generate_monthly_report():
    workflow = chain(
        fetch_sales_data.s(),
        calculate_metrics.s(),
        create_charts.s(),
        generate_pdf.s(),
        send_report_email.s()
    )
    return workflow.apply_async()

[OK] 17. DÉPLOIEMENT ET PRODUCTION


# Fichier requirements.txt
"""
celery[redis]==5.3.4
redis==5.0.1
flower==2.0.1
"""

# Systemd service (worker)
"""
[Unit]
Description=Celery Worker
After=network.target redis.service

[Service]
Type=forking
User=celery
Group=celery
WorkingDirectory=/app
Environment="CELERY_BROKER_URL=redis://localhost:6379/0"
ExecStart=/usr/local/bin/celery -A celery_app worker \
    --loglevel=info \
    --concurrency=4 \
    --max-tasks-per-child=100

[Install]
WantedBy=multi-user.target
"""

# Systemd service (beat)
"""
[Unit]
Description=Celery Beat
After=network.target redis.service

[Service]
Type=simple
User=celery
Group=celery
WorkingDirectory=/app
ExecStart=/usr/local/bin/celery -A celery_app beat --loglevel=info

[Install]
WantedBy=multi-user.target
"""

# Supervisord config
"""
[program:celery]
command=/path/to/venv/bin/celery -A celery_app worker --loglevel=info
directory=/app
user=celery
autostart=true
autorestart=true
stdout_logfile=/var/log/celery/worker.log
stderr_logfile=/var/log/celery/worker_err.log
"""

# Docker Compose
"""
version: '3.8'
services:
  redis:
    image: redis:7-alpine
    ports:
      - "6379:6379"
  
  worker:
    build: .
    command: celery -A celery_app worker --loglevel=info
    depends_on:
      - redis
    environment:
      - CELERY_BROKER_URL=redis://redis:6379/0
      - CELERY_RESULT_BACKEND=redis://redis:6379/0
  
  beat:
    build: .
    command: celery -A celery_app beat --loglevel=info
    depends_on:
      - redis
  
  flower:
    build: .
    command: celery -A celery_app flower
    ports:
      - "5555:5555"
    depends_on:
      - redis
"""

[OK] 18. BONNES PRATIQUES


# 1. Toujours définir timeout
@app.task(time_limit=300, soft_time_limit=250)
def long_running_task():
    pass

# 2. Idempotence (même résultat si exécuté plusieurs fois)
@app.task
def update_user_score(user_id, score):
    # Utiliser SET plutôt que INCREMENT
    User.objects.filter(id=user_id).update(score=score)

# 3. Éviter les tâches trop grandes
# Mauvais:
@app.task
def process_all_users():
    for user in User.objects.all():  # Millions d'utilisateurs!
        process_user(user)

# Bon:
@app.task
def process_user_batch(user_ids):
    for user_id in user_ids:
        process_user.delay(user_id)

# 4. Logging approprié
import logging
logger = logging.getLogger(__name__)

@app.task(bind=True)
def my_task(self):
    logger.info(f"Début tâche {self.request.id}")
    # ...
    logger.info(f"Fin tâche {self.request.id}")

# 5. Gestion des connexions DB
@app.task
def db_task():
    try:
        # Faire le travail
        pass
    finally:
        # Fermer connexions explicitement
        from django.db import connection
        connection.close()

[OK] FIN DU CHEATSHEET
