Celery en Python: los modos de fallo que un tutorial no te muestra
En 2025 el tutorial de Celery era así: instalas Redis, decoras una función con @app.task, arrancas un worker y todo funciona a la primera. Ese montaje sigue funcionando igual hoy, con Celery 5.6 en producción, pero deja fuera tres comportamientos que solo aparecen cuando un worker se reinicia a mitad de una tarea real: tareas que desaparecen sin dejar rastro, tareas que se ejecutan dos veces, y tareas envenenadas que, si aplicas al pie de la letra la configuración recomendada para arreglar las dos anteriores, se reencolan para siempre. Ninguno de los tres es un bug de Celery. Los tres nacen de decisiones de diseño explícitas sobre cuándo confirma Celery que un mensaje "ya se procesó", documentadas con sus propias advertencias, que la mayoría de configuraciones copiadas de un tutorial ignora.
Modo de fallo 1: la tarea que desaparece sin ejecutarse
Despliegas código nuevo, el worker se reinicia mientras procesaba una tarea de generación de informes, y esa tarea no vuelve a ejecutarse nunca. No hay excepción ni log de error.
La causa. Por defecto, Celery confirma (hace ack de) el mensaje antes de ejecutar la tarea, no después. Es una decisión deliberada: como el worker no puede saber si tu función es idempotente, Celery asume que no lo es y prefiere arriesgarse a perder una tarea antes que ejecutarla dos veces. La documentación oficial de tareas lo justifica así: el ack anticipado evita que una invocación que ya empezó se ejecute de nuevo si el worker muere a mitad de camino.
Cómo se detecta. No hay excepción que grepear: el síntoma es la ausencia de resultado. Si tienes tareas que "a veces no pasan nada" tras un despliegue o un OOM-kill del worker, y el log no muestra ni un intento, este es el sospechoso número uno antes de mirar cualquier otra cosa.
La mitigación. Si tu tarea es idempotente (repetirla con los mismos argumentos no rompe nada), task_acks_late = True mueve el ack a después de que la tarea termine. Esto arregla la pérdida silenciosa, pero abre la puerta al modo de fallo 2.
Modo de fallo 2: la misma tarea ejecutada dos veces
Activas acks_late para arreglar la pérdida silenciosa, y una tarea de 90 minutos termina ejecutándose dos veces en paralelo porque el broker la dio por perdida y la reenvió a otro worker mientras la primera seguía corriendo.
La causa. Con task_acks_late=True, si el worker muere durante la ejecución, el mensaje nunca llegó a confirmarse y el broker lo reenvía. Celery expone task_reject_on_worker_lost (deshabilitado por defecto) para forzar ese reencolado de forma explícita cuando el worker es señalado o muere abruptamente. Sobre Redis hay una tercera pieza: el visibility timeout, que por defecto son 3600 segundos (1 hora), el tiempo que Redis espera un ack antes de reenviar el mensaje a otro worker, exista o no acks_late de por medio.
Cómo se detecta. celery -A tasks inspect active y reserved: si ves la misma tarea con IDs de ejecución distintos corriendo casi al mismo tiempo, es duplicación por redelivery, no un bug en tu código de negocio.
La mitigación real no es un flag, es tu código. La idempotencia no se consigue con un comentario diciendo "esto es idempotente": se consigue con una restricción que el propio almacenamiento hace cumplir. Así se ve en la práctica, con una clave de idempotencia y una restricción única en base de datos:
CREATE TABLE informes_procesados (
idempotency_key TEXT PRIMARY KEY,
informe_id INTEGER NOT NULL,
procesado_en TIMESTAMPTZ NOT NULL DEFAULT now()
);
@app.task(bind=True, acks_late=True)
def generar_informe(self, informe_id: int, idempotency_key: str):
with db.begin():
fila = db.execute(
"SELECT 1 FROM informes_procesados WHERE idempotency_key = %s",
(idempotency_key,),
).fetchone()
if fila:
return {"informe_id": informe_id, "status": "ya_procesado"}
# La restricción UNIQUE de la tabla corta el duplicado si dos workers
# llegan aquí casi a la vez: el segundo INSERT no hace nada.
db.execute(
"INSERT INTO informes_procesados (idempotency_key, informe_id) "
"VALUES (%s, %s) ON CONFLICT (idempotency_key) DO NOTHING",
(idempotency_key, informe_id),
)
resultado = generar_contenido_informe(informe_id)
return {"informe_id": informe_id, "status": "completado", "resultado": resultado}
Si dos workers ejecutan la misma tarea casi a la vez tras un redelivery, el segundo INSERT falla por violación de unicidad o simplemente no inserta nada, y solo uno de los dos avanza a generar el informe (generar_contenido_informe es tu lógica de negocio real). Sin esa restricción en la base de datos, "idempotente" es solo una palabra en un comentario. Además, bajar worker_prefetch_multiplier a 1 (4 por defecto) reduce cuántas tareas quedan "en el aire" —ya sacadas de la cola pero no ejecutadas— cuando un proceso muere, lo que reduce la ventana en la que este modo de fallo aparece.
Modo de fallo 3: la tarea envenenada que se reencola para siempre
Este es el más grave de los tres, y el que la configuración recomendada en la mayoría de checklists de "buenas prácticas" deja sin resolver.
La causa. task_time_limit mata con SIGKILL el proceso worker que exceda el límite duro de tiempo y lo sustituye por uno nuevo. Si además tienes task_acks_late=True y task_reject_on_worker_lost=True —la combinación que cualquier checklist recomienda para no perder tareas— ese SIGKILL cuenta como worker perdido, y el mensaje se reencola. Si la tarea es venenosa, es decir, si siempre supera el time limit con esos argumentos, el ciclo se repite: se reencola, un worker la recoge, vuelve a superar el límite, muere, se reencola otra vez. La propia documentación de configuración de Celery lo advierte de forma explícita: activar task_reject_on_worker_lost "puede causar bucles de mensajes"; y sobre el semipredicado Reject con reencolado, la guía de tareas añade que reencolar con él "puede fácilmente resultar en un bucle de mensajes infinito" si no se usa con cuidado.
El motivo por el que max_retries de la tarea no te salva aquí es que este reencolado no pasa por self.retry(): es el broker redistribuyendo el mismo mensaje original tras dar por perdido al worker, no una llamada explícita de reintento contada por Celery. Celery no documenta ningún límite propio para esta ruta de reencolado por worker perdido, así que si no lo cortas tú, no lo corta nadie.
Cómo se detecta. El mismo task_id reapareciendo en celery -A tasks inspect active una y otra vez, con una duración que se acerca sistemáticamente a tu task_time_limit, sin que el contador de reintentos de tu propio código suba. Con Flower verás el mismo ID entrando y saliendo del estado "started" en bucle.
La mitigación. Cuenta tú mismo los reencolados por worker perdido —no puedes confiar en max_retries para esto— y usa Reject(requeue=False) para poner la tarea en cuarentena en vez de dejar que el broker la reencole otra vez:
from celery.exceptions import Reject
POISON_THRESHOLD = 3
@app.task(bind=True, acks_late=True)
def generar_informe_con_cuarentena(self, informe_id: int, idempotency_key: str):
# redis_client: tu cliente Redis (o el backend de resultados) para contar
# cuántas veces se ha reencolado este mismo mensaje.
contador_key = f"celery:worker_lost_count:{self.request.id}"
intentos = redis_client.incr(contador_key)
redis_client.expire(contador_key, 3600)
if intentos > POISON_THRESHOLD:
# registrar_en_cuarentena: tu tabla de dead-letter para inspección
# manual (o una dead-letter exchange si usas RabbitMQ).
registrar_en_cuarentena(
task_id=self.request.id,
informe_id=informe_id,
motivo="excede reencolados por worker perdido",
intentos=intentos,
)
raise Reject("tarea envenenada: en cuarentena para inspección manual", requeue=False)
return generar_informe(self, informe_id, idempotency_key)
A partir de POISON_THRESHOLD intentos, la tarea deja de reencolarse: queda registrada para inspección manual y el worker sigue disponible para el resto de la cola en vez de gastar ciclos repitiendo una tarea que nunca va a terminar bien.
Modo de fallo 4: el ETA que llega tarde o llega dos veces
Redis y RabbitMQ son, según la documentación de introducción de Celery, los dos únicos transportes de broker "feature complete"; el resto (incluida Amazon SQS) se documenta como experimental. Sobre Redis, el mismo visibility timeout del modo de fallo 2 tiene una segunda consecuencia que el tutorial básico nunca menciona.
La causa. Una tarea con countdown, eta o en su segundo o tercer retry cuyo tiempo total de espera supera el visibility timeout hace que Redis la dé por perdida y la reenvíe. Si la original sigue viva cuando termina, tienes dos ejecuciones simultáneas; si el patrón se repite en cada retry, un bucle equivalente al del modo de fallo 3, pero causado por el broker en vez de por el worker.
Cómo se detecta. Compara la duración real de tus tareas con ETA o countdown (percentil alto, no la media) contra tu visibility_timeout efectivo. Si alguna puede superarlo, tienes este modo de fallo esperando a pasar.
La mitigación. La documentación del broker Redis es explícita: subir el timeout para cubrir el ETA más largo "no es recomendable", porque solo retrasa la recuperación real de tareas perdidas por corte de energía o worker matado a la fuerza. Su alternativa para agendar a futuro lejano es otra: tareas periódicas respaldadas en base de datos, no un ETA sobre el broker. Si necesitas subirlo de todos modos, Celery exige tres ajustes coherentes a la vez, no uno solo:
app.conf.broker_transport_options = {'visibility_timeout': 43200}
app.conf.result_backend_transport_options = {'visibility_timeout': 43200}
app.conf.visibility_timeout = 43200 # 12 horas, en segundos
La configuración resultante, y lo que no cubre por sí sola
Punto de partida razonable para un worker con tareas idempotentes de duración variable (minutos, no segundos) sobre Redis, reuniendo las mitigaciones de los cuatro modos de fallo anteriores:
import os
from celery import Celery
app = Celery(
"myapp",
broker=os.environ["REDIS_BROKER_URL"],
backend=os.environ["REDIS_RESULT_BACKEND"],
)
app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
# Ack después de ejecutar: solo si tus tareas son idempotentes (modo 2).
task_acks_late=True,
# Reencola si el worker muere a mitad de tarea, en vez de perder el
# mensaje en silencio (modo 1). Sin el patrón de cuarentena del modo 3,
# esto por sí solo puede convertirse en un bucle de mensajes.
task_reject_on_worker_lost=True,
# Menos tareas reservadas por adelantado: menos trabajo "en el aire" si
# un proceso muere (modo 2).
worker_prefetch_multiplier=1,
# Límite duro de tiempo por tarea (soft = excepción capturable, hard =
# SIGKILL). Combinado con task_reject_on_worker_lost, esto es justo lo
# que dispara el modo de fallo 3 si la tarea es venenosa.
task_soft_time_limit=1800,
task_time_limit=1900,
# Mismo valor en los tres sitios: ver caveat del visibility timeout (modo 4).
broker_transport_options={"visibility_timeout": 3600},
result_backend_transport_options={"visibility_timeout": 3600},
visibility_timeout=3600,
result_expires=86400,
)
Esta lista de flags resuelve los modos de fallo 1, 2 y 4. El modo de fallo 3 no se resuelve con un flag: necesita el contador externo y el Reject(requeue=False) de la sección anterior, porque task_reject_on_worker_lost no viene con un límite de reencolados incorporado. Presentar esta lista como la configuración definitiva sin ese contador sería repetir exactamente el problema que abre este artículo: una tarea venenosa reencolándose para siempre.
Cómo reproducir estos fallos con un worker real
Para observar cualquiera de los diagnósticos de arriba hace falta primero una tarea corriendo. Con el app definido más arriba, en el mismo módulo o en un tasks.py que lo importe, levanta un worker:
celery -A tasks worker --loglevel=infoY encola una tarea desde una shell de Python o desde el endpoint que la dispare en tu API:
from tasks import generar_informe
import uuid
resultado = generar_informe.delay(42, str(uuid.uuid4()))
print(resultado.id, resultado.status) # PENDING: cubre "en cola" y "ejecutándose ahora"
Ese PENDING no significa "esperando a que el worker la recoja", como sugeriría la intuición. La documentación de configuración es explícita: sin activar task_track_started (deshabilitado por defecto), una tarea está pendiente, terminada, o esperando un reintento, sin un estado intermedio de "en ejecución". Si necesitas ver ese matiz, por ejemplo para un dashboard de progreso, activa task_track_started=True y el estado pasará a STARTED en cuanto un worker la tome.
Con esto ya puedes reproducir celery -A tasks inspect active/reserved y observar el comportamiento de ack descrito arriba. Si lo que necesitas no es lanzar una tarea puntual sino repetirla cada hora, no la dispares con un cron externo golpeando tu API: usa Celery Beat, el scheduler integrado para tareas periódicas, y aplícale la misma disciplina de idempotencia que al resto, porque un reinicio de Beat o un reloj desincronizado puede disparar la misma tarea programada dos veces.
Cuándo Celery es la pieza correcta y cuándo es demasiada pieza
Celery no es la única respuesta a "necesito ejecutar esto en segundo plano", y para bastantes casos es la opción más pesada de instalar y operar. La propia documentación de FastAPI lo dice sin rodeos: usa su BackgroundTasks integrado para tareas pequeñas dentro del mismo proceso, y reserva herramientas como Celery, con RabbitMQ o Redis detrás, para trabajo que necesita correr en varios procesos o varios servidores sin compartir memoria.
| Herramienta | Cuándo tiene sentido | Cuándo no |
|---|---|---|
| FastAPI BackgroundTasks | Tareas cortas en el mismo proceso, sin reintentos ni cola persistente (ej. enviar un email tras responder) | Cualquier cosa que deba sobrevivir un reinicio del proceso o escalar a varios servidores |
| Celery | Pipelines de datos o IA, tareas periódicas (Beat), flujos con chains/chords/groups, necesitas RabbitMQ además de Redis | Un equipo pequeño sin experiencia operando colas: el coste de infraestructura y depuración es real |
| RQ | Solo Redis o Valkey, cola simple, prioriza baja barrera de entrada sobre features | Necesitas AMQP/RabbitMQ, o workflows complejos tipo chord |
| arq | Stack ya async de punta a punta (FastAPI/asyncio), Redis, mantenimiento activo | Necesitas un broker distinto de Redis, o el ecosistema de monitorización de Celery (Flower) |
| Dramatiq | Quieres algo más simple que Celery pero con soporte real de RabbitMQ y Redis | Dependes de piezas específicas del ecosistema Celery, como Beat o Canvas avanzado |
La pregunta que de verdad decide no es qué librería es mejor en abstracto, sino si tu tarea es idempotente y qué pasa el día en que un worker muere a mitad de ejecutarla, o en que una tarea concreta nunca puede terminar a tiempo. Si todavía no puedes responder eso, ese es el problema que hay que resolver antes de elegir broker, backend o librería.