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.
Incluida en el Pase · para Claude Code, Cursor, Codex CLI
""" 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
// opiniones_de_la_comunidad
Opiniones
Cargando opiniones…