Celery es una cola de tareas distribuida que ejecuta trabajo fuera del ciclo de una petición. En lugar de hacer que el usuario espere un correo, un informe o el procesamiento de una imagen, la aplicación publica un mensaje y un worker ejecuta la función. La decisión directa es esta: usa Celery cuando el trabajo pueda terminar después, necesite reintentos o deba distribuirse entre procesos y máquinas. No lo agregues solo para invocar una función breve.
Celery coordina productores y workers mediante un broker; no es el broker ni la base de datos. En este ejemplo, Redis transportará los mensajes y podrá guardar resultados si son necesarios. Consulta nuestra guía de Redis con Python y utiliza la documentación oficial de Celery para los detalles de cada versión.
Arquitectura y decisiones iniciales
Una vista HTTP valida datos, confirma el estado necesario y envía una tarea con .delay() o .apply_async(). El broker conserva el mensaje hasta que un worker lo reserva. Un backend opcional registra el estado y el retorno. Esto reduce la latencia de la API, pero introduce consistencia eventual y fallos distribuidos. Devuelve un identificador y ofrece consulta, webhook o eventos si el cliente necesita conocer el progreso.
Celery encaja en trabajos de segundos o minutos, tareas programadas, colas aisladas y procesamiento recuperable. Para concurrencia de I/O dentro de un proceso, async/await en Python resuelve otro problema y no vuelve durable una operación. Los mecanismos de background del framework son suficientes solo para tareas pequeñas que se pueden perder junto al proceso web.
Configuración con Redis
Instala el extra de Redis y ejecuta Redis localmente o como servicio administrado. En producción, configura autenticación, TLS cuando corresponda, persistencia según la tolerancia a pérdidas y alertas de memoria. Nunca lo expongas directamente a internet. La documentación oficial de Redis explica estas opciones.
pip install "celery[redis]"
celery -A tasks worker --loglevel=INFO
from celery import Celery
app = Celery(
"myapp",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
task_track_started=True,
task_time_limit=300,
task_soft_time_limit=270,
worker_prefetch_multiplier=1,
)
Carga las URLs desde variables de entorno. JSON evita deserializar objetos Python arbitrarios. El backend de resultados consume memoria y necesita caducidad; usa task_ignore_result si nadie leerá el retorno. Para empaquetar web y workers de forma repetible, revisa Docker con Python.
Reintentos que favorecen la recuperación
Reintenta fallos transitorios: timeouts, indisponibilidad temporal y límites de solicitudes. Datos inválidos o reglas de negocio rechazadas son errores permanentes; repetirlos bloquea trabajo útil. El backoff exponencial, el jitter y un máximo de intentos evitan que todos los workers ataquen al mismo tiempo un servicio que se recupera.
import requests
@app.task(bind=True, max_retries=5)
def notify_customer(self, order_id: int) -> None:
try:
response = requests.post(
"https://api.example.com/notifications",
json={"order_id": order_id},
timeout=(3, 10),
)
except (requests.Timeout, requests.ConnectionError) as error:
espera = min(2 ** self.request.retries, 600)
raise self.retry(exc=error, countdown=espera)
if response.status_code in {429, 500, 502, 503, 504}:
retry_after = response.headers.get("Retry-After")
espera = int(retry_after) if retry_after and retry_after.isdigit() else min(2 ** self.request.retries, 600)
raise self.retry(exc=requests.HTTPError(response=response), countdown=espera)
response.raise_for_status() # Los 4xx permanentes fallan sin reintento.</code></pre>
Define timeouts de red aunque la tarea tenga límite de tiempo. No captures cualquier Exception para repetir sin criterio. Registra la causa y conserva la excepción original. La guía oficial de reintentos de Celery detalla retry y autoretry_for.
La idempotencia es obligatoria
Las colas suelen ofrecer entrega al menos una vez. Un worker puede completar un efecto externo, morir antes de confirmar y recibir el mismo mensaje. Cobrar una tarjeta, emitir un cupón o incrementar un saldo debe tolerar esa repetición.
Genera una clave estable como charge:order:123, impón unicidad en la base de datos y actualiza estado y efecto en una transacción cuando sea posible. Si una API externa admite claves de idempotencia, envía siempre la misma. Un lock en Redis reduce concurrencia, pero no sustituye una restricción única durable.
@app.task(bind=True, acks_late=True)
def issue_invoice(self, order_id: int) -> str:
invoice = Invoice.get_or_create_for_order(order_id)
if invoice.status == "issued":
return invoice.number
number = tax_api.issue(
order_id=order_id,
idempotency_key=f"invoice-order-{order_id}",
)
invoice.mark_issued(number)
return number</code></pre>
acks_late=True confirma después de ejecutar y permite recuperar trabajo si el proceso muere, aunque hace más probable una repetición. Úsalo solo con tareas idempotentes. Envía IDs y valores simples, no objetos ORM serializados que podrían quedar obsoletos.
Colas, concurrencia y presión
Separa cargas por comportamiento, por ejemplo emails, reports y payments. Un informe pesado no bloqueará confirmaciones de pago. Enruta con task_routes y reserva workers. Las tareas de CPU requieren procesos; las de I/O admiten mayor concurrencia solo si la base y los servicios remotos la toleran. Nuestra comparación de threads y procesos en Python ayuda a distinguirlas.
Un prefetch alto mejora trabajos diminutos, pero permite que un worker reserve demasiadas tareas largas. Un multiplicador de uno es un punto inicial razonable. Aplica contrapresión: limita envíos, agrupa con cuidado o degrada trabajo opcional antes de agotar la memoria del broker.
Observabilidad en producción
Mide profundidad y edad de la cola, duración, throughput, reintentos, fallos definitivos y disponibilidad de workers. Los logs estructurados deben contener ID, nombre, intento, cola e identificador de negocio, nunca secretos. Propaga un correlation ID para conectar la petición web con la ejecución.
Los eventos de Celery y Flower ayudan a inspeccionar, pero no sustituyen métricas históricas y alertas. Combina logging estructurado en Python con trazas y un proceso para tareas agotadas. Alerta por impacto, como la tarea más antigua superando el objetivo de servicio, no por la mera existencia de mensajes.
Errores frecuentes
- Publicar antes de confirmar la transacción, por lo que el worker no encuentra el registro.
- Llamar servicios sin timeout o reintentar errores permanentes indefinidamente.
- Suponer entrega exactamente una vez y omitir la idempotencia.
- Enviar archivos grandes en el mensaje en vez de guardar y pasar una referencia.
- Mezclar tareas rápidas y pesadas sin colas ni límites separados.
- Cambiar una firma incompatible mientras quedan mensajes antiguos.
Publica después del commit. Si no puedes tolerar el hueco entre base de datos y broker, utiliza un outbox transaccional. Versiona mensajes y conserva compatibilidad mientras se vacía la cola. Prueba la función de dominio sin Celery y añade integración para publicación, reintento y duplicados; la guía de pytest aporta una base práctica.
Preguntas frecuentes
¿Celery sustituye a async/await?
No. Celery distribuye trabajo y puede persistir mensajes; async/await coordina operaciones concurrentes dentro de un proceso. Ambos pueden convivir.
¿Redis es suficiente como broker de producción?
Puede serlo si durabilidad, disponibilidad, memoria y seguridad coinciden con el riesgo. Considera RabbitMQ cuando necesites enrutamiento o semánticas avanzadas.
¿Debo guardar todos los resultados?
No. Guárdalos solo si un consumidor los consulta. Para efectos laterales, el estado de negocio en la base suele ser la fuente de verdad.
¿Cómo pruebo reintentos sin esperar?
Prueba el dominio por separado, simula excepciones transitorias concretas y verifica la solicitud de retry. Completa con pocos tests usando broker y worker reales.
Conclusión
Una configuración confiable de Celery va mucho más allá de .delay(). Clasifica errores, limita reintentos, vuelve idempotentes los efectos, aísla cargas, controla concurrencia y observa tanto espera como ejecución. Empieza con una tarea acotada, Redis protegido, mensajes pequeños y un objetivo medible. Escala workers y rutas solo cuando esas garantías estén claras.