Cómo construir un pipeline de procesamiento por lotes en Python con FastAPI, Celery y Redis

python Cómo construir un pipeline de procesamiento por lotes en Python con FastAPI, Celery y Redis

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.

Comentarios
¿Quieres comentar?

Inicia sesión con Telegram para participar en la conversación


Comentarios (0)

Aún no hay comentarios. ¡Sé el primero en comentar!

Iniciar Sesión