← Volver al catálogo
IA, ingeniería y MLOpsReferenciaIntermedioEn el pase

Tareas Asíncronas en Django con Celery

Patrones de nivel producción para procesar tareas en segundo plano en Django usando Celery con Redis o RabbitMQ. Cubre configuración, diseño de tareas, planificación con Celery Beat, reintentos, flujos canvas, monitorización y testing. Útil cuando hay que descargar el ciclo de petición de operaciones lentas (correos, generación de PDF, llamadas a API) o programar tareas periódicas.

Comprobando acceso…

Incluida en el Pase · para Claude Code, Cursor, Codex CLI

// resultado_de_ejemplo

""" LeadPulse SaaS — Tareas Asíncronas con Django + Celery

Implementación de producción basada en el skill tareas-asincronas-django-celery.

Estructura de archivos: leadpulse/ ├── config/ │ ├── init.py │ ├── celery.py │ └── settings/ │ ├── base.py (fragmento relevante) │ └── test.py ├── apps/ │ ├── leads/ │ │ └── tasks.py ← tareas críticas │ ├── billing/ │ │ └── tasks.py ← cobro idempotente │ ├── reports/ │ │ └── tasks.py ← pipeline canvas │ └── core/ │ └── signals.py ← integración Sentry └── tests/ └── test_lead_tasks.py """

═══════════════════════════════════════════════════════════════════

config/celery.py — Entrypoint de la app Celery

═══════════════════════════════════════════════════════════════════

import os from celery import Celery

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

app = Celery('leadpulse') app.config_from_object('django.conf:settings', namespace='CELERY') app.autodiscover_tasks()

═══════════════════════════════════════════════════════════════════

config/init.py

═══════════════════════════════════════════════════════════════════

from .celery import app as celery_app # noqa: F401, E402

all = ('celery_app',)

═══════════════════════════════════════════════════════════════════

config/settings/base.py — Fragmento de configuración Celery

═══════════════════════════════════════════════════════════════════

from celery.schedules import crontab # noqa: E402

-- Broker y backend --------------------------------------------------

CELERY_BROKER_URL = 'redis://localhost:6379/0' # Railway: var de entorno CELERY_RESULT_BACKEND = 'django-db' # django-celery-results

-- Serialización ------------------------------------------------------

CELERY_ACCEPT_CONTENT = ['json'] CELERY_TASK_SERIALIZER = 'json' CELERY_RESULT_SERIALIZER = 'json'

-- Comportamiento de tareas -------------------------------------------

CELERY_TASK_TRACK_STARTED = True CELERY_TASK_TIME_LIMIT = 30 * 60 # Hard: 30 min CELERY_TASK_SOFT_TIME_LIMIT = 25 * 60 # Soft: 25 min → limpieza antes del kill CELERY_WORKER_PREFETCH_MULTIPLIER = 1 # Fair scheduling de tareas largas CELERY_TASK_ACKS_LATE = True # Re-encolar si el worker cae

-- Resultados ---------------------------------------------------------

CELERY_RESULT_EXPIRES = 60 * 60 * 24 # Guardar resultados 24 h

-- Beat scheduler (DB) ------------------------------------------------

CELERY_BEAT_SCHEDULER = 'django_celery_beat.schedulers:DatabaseScheduler'

-- Tareas periódicas definidas en código (fallback) ------------------

CELERY_BEAT_SCHEDULE = { 'cleanup-inactive-leads': { 'task': 'leads.cleanup_inactive_leads', 'schedule': crontab(hour=3, minute=0), # 3am cada noche }, 'sync-email-templates': { 'task': 'leads.sync_email_templates', 'schedule': 3600.0, # cada hora }, 'weekly-kpi-report': { 'task': 'reports.send_weekly_kpi_report', 'schedule': crontab(day_of_week='monday', hour=7, minute=0), }, }

INSTALLED_APPS = [ # ... apps de Django ... 'django_celery_results', 'django_celery_beat', ]

═══════════════════════════════════════════════════════════════════

apps/leads/tasks.py — Tareas críticas de leads

═══════════════════════════════════════════════════════════════════

import logging from celery import shared_task from celery.exceptions import SoftTimeLimitExceeded

logger = logging.getLogger(name)

@shared_task(name='leads.send_lead_email') def send_lead_email(lead_id: int, template: str = 'welcome') -> None: """ Envía email de bienvenida/nurturing a un nuevo lead. Idempotente: si el lead no existe, termina silenciosamente. Cola: high_priority """ from apps.leads.models import Lead from apps.leads.services import EmailService

try:
    lead = Lead.objects.get(pk=lead_id)
except Lead.DoesNotExist:
    logger.warning('send_lead_email: lead %s no encontrado, descartando', lead_id)
    return  # No lanzar excepción — tarea ya no completable

EmailService.send(lead=lead, template=template)
logger.info('Email "%s" enviado a lead %s (%s)', template, lead_id, lead.email)

@shared_task( bind=True, name='leads.enrich_lead_profile', max_retries=5, default_retry_delay=30, autoretry_for=(ConnectionError, TimeoutError), retry_backoff=True, retry_backoff_max=600, retry_jitter=True, ) def enrich_lead_profile(self, lead_id: int) -> dict: """ Llama a la API de Clearbit para enriquecer datos del lead. Reintentable con backoff exponencial ante fallos de red. Cola: high_priority """ from apps.leads.models import Lead from apps.leads.services import ClearbitClient

lead = Lead.objects.get(pk=lead_id)
try:
    data = ClearbitClient().enrich(email=lead.email)
    Lead.objects.filter(pk=lead_id).update(
        company=data.get('company'),
        job_title=data.get('title'),
        linkedin_url=data.get('linkedin'),
        enriched=True,
    )
    logger.info('Lead %s enriquecido correctamente', lead_id)
    return data
except ClearbitClient.RateLimitError as exc:
    raise self.retry(exc=exc, countdown=int(exc.retry_after))

@shared_task( name='leads.cleanup_inactive_leads', soft_time_limit=120, time_limit=150, ) def cleanup_inactive_leads() -> int: """ Archiva leads sin actividad en los últimos 90 días. Se ejecuta a las 3am vía Celery Beat. Retorna el número de leads archivados. """ from django.utils import timezone from datetime import timedelta from apps.leads.models import Lead

try:
    cutoff = timezone.now() - timedelta(days=90)
    updated = Lead.objects.filter(
        last_activity__lt=cutoff,
        status=Lead.Status.ACTIVE,
    ).update(status=Lead.Status.ARCHIVED)
    logger.info('cleanup_inactive_leads: %d leads archivados', updated)
    return updated
except SoftTimeLimitExceeded:
    logger.warning('cleanup_inactive_leads: tiempo límite alcanzado, limpiando')
    raise

@shared_task(name='leads.sync_email_templates') def sync_email_templates() -> None: """Sincroniza plantillas de email desde el servicio externo (cada hora).""" from apps.leads.services import TemplateSync TemplateSync.run() logger.info('sync_email_templates: plantillas actualizadas')

═══════════════════════════════════════════════════════════════════

apps/billing/tasks.py — Cobro idempotente vía Stripe

═══════════════════════════════════════════════════════════════════

@shared_task( bind=True, name='billing.charge_subscription', max_retries=3, ) def charge_subscription(self, subscription_id: int) -> None: """ Cobra la suscripción mensual del cliente vía Stripe. Guard idempotente: si ya está cobrada, termina sin error. Cola: high_priority """ from apps.billing.models import Subscription, FailedCharge

sub = Subscription.objects.select_for_update().get(pk=subscription_id)

# Guard idempotente — seguro ante reeintentos
if sub.status != Subscription.Status.PENDING_CHARGE:
    logger.info(
        'charge_subscription: suscripción %s ya procesada (estado: %s)',
        subscription_id, sub.status,
    )
    return

try:
    _do_stripe_charge(sub)
    sub.status = Subscription.Status.ACTIVE
    sub.save(update_fields=['status'])
except Exception as exc:
    if self.request.retries >= self.max_retries:
        # Dead-letter: persistir para revisión manual
        FailedCharge.objects.create(
            subscription_id=subscription_id,
            error=str(exc),
            task_id=self.request.id,
        )
        logger.error(
            'charge_subscription: máximo de reintentos alcanzado para suscripción %s',
            subscription_id,
        )
        return
    raise self.retry(exc=exc, countdown=60 * (2 ** self.request.retries))

def _do_stripe_charge(sub): """Lógica real de cobro (stub).""" import stripe stripe.PaymentIntent.create( amount=int(sub.amount_eur * 100), currency='eur', customer=sub.stripe_customer_id, confirm=True, )

═══════════════════════════════════════════════════════════════════

apps/reports/tasks.py — Pipeline canvas y reportes

═══════════════════════════════════════════════════════════════════

from celery import chain, group, chord, shared_task # noqa: E402

@shared_task(name='reports.fetch_leads_chunk') def fetch_leads_chunk(source_id: int) -> list: """Descarga chunk de leads desde CSV importado.""" from apps.leads.services import CSVImporter return CSVImporter.fetch_chunk(source_id)

@shared_task(name='reports.enrich_leads_batch') def enrich_leads_batch(lead_ids: list) -> list: """Enriquece un batch de leads en paralelo.""" return lead_ids # delegado a enrich_lead_profile individualmente

@shared_task(name='reports.score_leads') def score_leads(lead_data: list) -> list: """Calcula el score de cada lead según modelo ML.""" from apps.leads.services import LeadScorer return LeadScorer.score_batch(lead_data)

@shared_task(name='reports.notify_sales_team') def notify_sales_team(scored_leads: list) -> None: """Notifica al equipo comercial con los leads puntuados.""" from apps.notifications.services import SlackNotifier SlackNotifier.send_lead_batch(scored_leads)

def import_leads_pipeline(source_id: int, lead_ids: list) -> None: """ Pipeline completo de importación de leads vía Canvas: fetch → enrich (paralelo) → score → notificar al comercial """ pipeline = chain( fetch_leads_chunk.s(source_id), chord( group(enrich_lead_profile.s(lid) for lid in lead_ids), score_leads.s(), ), notify_sales_team.s(), ) pipeline.delay()

@shared_task(name='reports.send_weekly_kpi_report') def send_weekly_kpi_report() -> None: """ Genera y envía reporte semanal de KPIs a los admins. Se ejecuta los lunes a las 7am vía Celery Beat. """ from apps.reports.services import KPIReport from apps.users.models import User

report = KPIReport.generate()
admins = User.objects.filter(is_staff=True).values_list('email', flat=True)
for email in admins:
    send_lead_email.apply_async(
        args=[None],
        kwargs={'template': 'weekly_kpi', 'recipient': email, 'data': report},
        queue='default',
    )
logger.info('Reporte KPI semanal enviado a %d admins', len(admins))

═══════════════════════════════════════════════════════════════════

apps/core/signals.py — Integración Sentry para fallos de tareas

═══════════════════════════════════════════════════════════════════

from celery.signals import task_failure # noqa: E402

@task_failure.connect def on_task_failure(sender, task_id, exception, args, kwargs, traceback, einfo, **kw): """Captura todos los fallos de tareas Celery en 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, }) scope.set_tag('celery.task', sender.name) sentry_sdk.capture_exception(exception)

═══════════════════════════════════════════════════════════════════

tests/test_lead_tasks.py — Suite de tests

═══════════════════════════════════════════════════════════════════

import pytest from unittest.mock import patch, MagicMock

class TestSendLeadEmail:

@pytest.mark.django_db
def test_sends_email_to_existing_lead(self, lead_fixture):
    with patch('apps.leads.services.EmailService') as mock_email:
        send_lead_email(lead_fixture.pk, template='welcome')
        mock_email.send.assert_called_once_with(
            lead=lead_fixture, template='welcome'
        )

@pytest.mark.django_db
def test_skips_missing_lead_gracefully(self):
    """No debe lanzar si el lead fue borrado entre encolar y ejecutar."""
    send_lead_email(99999)  # Lead inexistente — debe terminar sin error

class TestEnrichLeadProfile:

@pytest.mark.django_db
def test_enriches_lead_data(self, lead_fixture):
    mock_data = {'company': 'Acme Inc.', 'title': 'CTO', 'linkedin': 'linkedin.com/in/acme'}
    with patch('apps.leads.services.ClearbitClient.enrich', return_value=mock_data):
        result = enrich_lead_profile(lead_fixture.pk)
        assert result['company'] == 'Acme Inc.'

@pytest.mark.django_db
def test_retries_on_connection_error(self, lead_fixture):
    with patch('apps.leads.services.ClearbitClient.enrich') as mock_enrich:
        mock_enrich.side_effect = ConnectionError('timeout')
        with pytest.raises(ConnectionError):
            enrich_lead_profile.apply(args=[lead_fixture.pk], throw=True)

class TestChargeSubscription:

@pytest.mark.django_db
def test_skips_already_charged_subscription(self, subscription_fixture):
    subscription_fixture.status = 'ACTIVE'
    subscription_fixture.save()
    with patch('apps.billing.tasks._do_stripe_charge') as mock_charge:
        charge_subscription(subscription_fixture.pk)
        mock_charge.assert_not_called()

@pytest.mark.django_db
def test_creates_failed_charge_after_max_retries(self, subscription_fixture):
    with patch('apps.billing.tasks._do_stripe_charge') as mock_charge:
        mock_charge.side_effect = Exception('Card declined')
        charge_subscription.apply(args=[subscription_fixture.pk])
        from apps.billing.models import FailedCharge
        assert FailedCharge.objects.filter(
            subscription_id=subscription_fixture.pk
        ).exists()

═══════════════════════════════════════════════════════════════════

config/settings/test.py — Settings para CI

═══════════════════════════════════════════════════════════════════

CELERY_TASK_ALWAYS_EAGER = True # Ejecutar tareas síncronamente en tests CELERY_TASK_EAGER_PROPAGATES = True # Re-lanzar excepciones desde tareas

═══════════════════════════════════════════════════════════════════

Comandos de arranque (README / Makefile)

═══════════════════════════════════════════════════════════════════

STARTUP_COMMANDS = """

Desarrollo (worker + beat combinados, NUNCA en producción)

celery -A config worker --beat --loglevel=info

Producción — workers separados por cola

celery -A config worker --loglevel=warning --concurrency=4 -Q high_priority,default,low_priority celery -A config beat --loglevel=info --scheduler django_celery_beat.schedulers:DatabaseScheduler

Monitorización visual con Flower

celery -A config flower --port=5555

Inspección rápida desde CLI

celery -A config inspect active celery -A config inspect stats redis-cli llen celery # longitud de la cola default """

═══════════════════════════════════════════════════════════════════

Checklist de producción

═══════════════════════════════════════════════════════════════════

PRODUCTION_CHECKLIST = { "CELERY_TASK_ACKS_LATE": True, # ✓ Re-encolar ante crash de worker "CELERY_WORKER_PREFETCH_MULTIPLIER": 1, # ✓ Distribución justa de tareas largas "CELERY_TASK_SOFT_TIME_LIMIT": 25 * 60, # ✓ Cleanup antes del hard kill "sentry_integration": True, # ✓ Captura task_failure "separate_queues": [ # ✓ Aislamiento por prioridad "high_priority", "default", "low_priority", ], "flower_monitoring": "port 5555", # ✓ Visibilidad en tiempo real "beat_single_node": True, # ✓ Un solo proceso Beat en prod "worker_supervisor": "Railway health", # ✓ Reinicio automático en crash }

// qué_hace

Implementa procesamiento de tareas en segundo plano y programadas en aplicaciones Django con Celery.

// cómo_lo_hace

Aporta patrones de configuración del app entrypoint y settings de Celery, diseño de tareas idempotentes con reintentos, planificación periódica con Celery Beat, composición de flujos canvas, monitorización de colas y estrategias de testing de tareas.

// ejemplo_de_uso

Úsala cuando una tarea pesada bloquea la respuesta de tu app Django y el usuario espera mirando una pantalla colgada. Ej.: mueves el envío de un informe PDF a una tarea Celery con reintentos y programas un resumen diario con Celery Beat a las 8:00.

// plataformas

Claude CodeCursorCodex CLI
Categoría
IA, ingeniería y MLOps
Tipo
Referencia
Nivel
Intermedio
Licencia
MIT
Seguridad
seguro · riesgo bajo
Versión
1.0.0

// opiniones_de_la_comunidad

Opiniones

Cargando opiniones…