Cómo construir un pipeline de procesamiento por lotes en Python con FastAPI, Celery y Redis
En este tutorial práctico vas a crear un servicio que reciba archivos vía HTTP, encole tareas de procesamiento con Celery (broker: Redis) y devuelva el estado y resultados. Incluye estructura de proyecto, código completo y explicación de por qué tomar estas decisiones.
¿Qué vas a aprender?
- Configurar FastAPI para subir archivos y encolar trabajos
- Configurar Celery con Redis como broker y backend
- Diseñar tareas idempotentes y manejo básico de errores
- Desplegar localmente con docker-compose
Estructura del proyecto
project/
├─ app/
│ ├─ main.py # FastAPI app
│ ├─ tasks.py # Definición de tareas Celery
│ ├─ worker.py # Inicializador del worker (opcional)
│ ├─ storage.py # Helpers para guardar/leer archivos
│ └─ config.py # Configuración compartida
├─ docker-compose.yml
├─ Dockerfile.app
├─ Dockerfile.worker
├─ requirements.txt
└─ README.md
Dependencias (requirements.txt)
fastapi==0.95.2
uvicorn[standard]==0.22.0
celery[redis]==5.3.1
redis==4.5.0
python-multipart==0.0.6
aiofiles==23.1.0
boto3==1.31.0 # opcional para S3
config.py
import os
REDIS_URL = os.getenv('REDIS_URL', 'redis://redis:6379/0')
RESULT_BACKEND = os.getenv('RESULT_BACKEND', 'redis://redis:6379/1')
UPLOAD_DIR = os.getenv('UPLOAD_DIR', '/data/uploads')
storage.py
import os
from pathlib import Path
from config import UPLOAD_DIR
Path(UPLOAD_DIR).mkdir(parents=True, exist_ok=True)
def save_upload(file_obj, filename: str) -> str:
path = Path(UPLOAD_DIR) / filename
with open(path, 'wb') as f:
while True:
chunk = file_obj.read(1024*1024)
if not chunk:
break
f.write(chunk)
return str(path)
def read_result(path: str) -> bytes:
with open(path, 'rb') as f:
return f.read()
tasks.py
from celery import Celery
from config import REDIS_URL, RESULT_BACKEND
import time
import os
from storage import save_upload
celery = Celery('tasks', broker=REDIS_URL, backend=RESULT_BACKEND)
@celery.task(bind=True, acks_late=True)
def process_file(self, filepath: str, output_name: str) -> dict:
"""
Tarea de ejemplo que 'procesa' un archivo. Aquí puede ir OCR, resize, análisis, etc.
acks_late=True para evitar perder tareas si el worker cae durante el procesamiento.
"""
try:
# Simular trabajo pesado
for i in range(5):
time.sleep(1)
self.update_state(state='PROGRESS', meta={'current': i+1, 'total': 5})
# Ejemplo: crear un archivo de salida que indique éxito
out_path = f"{filepath}.processed"
with open(out_path, 'w') as f:
f.write(f"processed {os.path.basename(filepath)}\n")
return {'status': 'ok', 'output': out_path}
except Exception as e:
raise self.retry(exc=e, countdown=5, max_retries=3)
main.py (FastAPI)
from fastapi import FastAPI, UploadFile, File, HTTPException
from fastapi.responses import JSONResponse, FileResponse
from tasks import process_file
from storage import save_upload
import uuid
from config import UPLOAD_DIR
import os
app = FastAPI()
@app.post('/upload')
async def upload(file: UploadFile = File(...)):
# Generar nombre único
uid = str(uuid.uuid4())
filename = f"{uid}_{file.filename}"
filepath = os.path.join(UPLOAD_DIR, filename)
# Guardar el archivo a disco de manera eficiente
with open(filepath, 'wb') as f:
content = await file.read()
f.write(content)
# Encolar tarea
task = process_file.apply_async(args=[filepath, filename])
return JSONResponse({'task_id': task.id, 'status_url': f"/status/{task.id}"})
@app.get('/status/{task_id}')
def status(task_id: str):
res = process_file.AsyncResult(task_id)
info = {
'task_id': task_id,
'state': res.state,
'info': res.info
}
return info
@app.get('/result/{task_id}')
def result(task_id: str):
res = process_file.AsyncResult(task_id)
if res.state != 'SUCCESS':
raise HTTPException(status_code=404, detail='Result not ready')
output_path = res.result.get('output')
return FileResponse(output_path, media_type='application/octet-stream', filename=os.path.basename(output_path))
Docker Compose (docker-compose.yml)
version: '3.8'
services:
redis:
image: redis:7
ports:
- '6379:6379'
web:
build:
context: .
dockerfile: Dockerfile.app
volumes:
- ./data:/data
ports:
- '8000:8000'
environment:
- REDIS_URL=redis://redis:6379/0
- RESULT_BACKEND=redis://redis:6379/1
- UPLOAD_DIR=/data/uploads
depends_on:
- redis
worker:
build:
context: .
dockerfile: Dockerfile.worker
volumes:
- ./data:/data
environment:
- REDIS_URL=redis://redis:6379/0
- RESULT_BACKEND=redis://redis:6379/1
- UPLOAD_DIR=/data/uploads
depends_on:
- redis
Dockerfiles
# Dockerfile.app
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY app /app
CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]
# Dockerfile.worker
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY app /app
CMD ["celery", "-A", "tasks.celery", "worker", "--loglevel=info", "--concurrency=2"]
Arrancar el entorno
docker-compose up --build
Con esto tendrás Redis, la API y un worker funcionando localmente.
Probar el flujo
# Subir un archivo
curl -F "file=@./sample.txt" http://localhost:8000/upload
# Comprobar estado (usar el task_id devuelto)
curl http://localhost:8000/status/
# Descargar resultado cuando esté listo
curl -OJ http://localhost:8000/result/
Por qué esta arquitectura
- FastAPI gestiona peticiones HTTP de forma asíncrona y eficiente.
- Celery desacopla trabajo en segundo plano, permitiendo escalado horizontal de workers.
- Redis como broker y backend es rápido, sencillo de desplegar y suficiente para muchos casos de uso.
- Guardar archivos en disco (o S3) y pasar rutas a la tarea evita mover blobs por el broker, reduciendo latencia y uso de memoria.
Buenas prácticas y consideraciones
- Usa acks_late y retries para resiliencia.
- Haz las tareas idempotentes: si se procesan dos veces no deben corromper estado.
- No pases archivos binarios grandes por el broker; pasa rutas o referencias (S3 presigned URLs).
- Controla el almacenamiento temporal: limpia /data periódicamente.
- Configura límites de concurrencia y timeouts para evitar sobrecargar workers.
Escalado y monitorización
Para producción considera:
- Usar S3/MinIO para almacenamiento en lugar del filesystem local.
- Configurar Flower para monitorizar tareas Celery.
- Separar Redis para broker y backend (o usar backend persistente como DB si necesitas historial).
- Agregar autenticación y límites por usuario (rate limiting, cuotas).
Consejo avanzado: habilita tracing (OpenTelemetry) tanto en FastAPI como en Celery para correlacionar peticiones y tareas en logs y trazas distribuidas; facilita depuración en sistemas distribuidos.
Advertencia: en producción, valida y escanea uploads (tamaño, tipo), protege endpoints y no confíes en nombres de archivos aceptados tal cual; evita path traversal al guardar archivos.
Siguiente paso: integra S3 para almacenamiento y añade un job que limpie archivos temporales y registre métricas en Prometheus para tomar decisiones de autoscaling.
¿Quieres comentar?
Inicia sesión con Telegram para participar en la conversación