← Volver al catálogo
IA, ingeniería y MLOpsReferenciaIntermedioGratis

Trabajos en Segundo Plano con Python

Patrones para implementar colas de tareas, workers y arquitecturas orientadas a eventos en Python. Ideal para desacoplar operaciones largas del ciclo peticion-respuesta usando Celery, RQ u otras soluciones.

Descargar SKILL.md

Descarga abierta · sin registro · para Python, Celery, Redis

// resultado_de_ejemplo

""" CULTIVA IA — Sistema de Trabajos en Segundo Plano

Stack: Python 3.11 · FastAPI · Celery 5.x · Redis · PostgreSQL (SQLAlchemy async)

Este módulo implementa los patrones de la skill trabajos-en-segundo-plano-python aplicados a la plataforma real de CULTIVA IA:

• Exportación de reportes de campañas (CSV/PDF) • Generación de contenido AI en lote (LLM) • Envío de webhooks a integraciones externas • Procesamiento de catálogos de producto (CSV upload)

Arquitectura ──────────── FastAPI ──► Redis (broker) ──► Celery Workers │ ▼ PostgreSQL (jobs table) │ ▼ S3 / GCS (artefactos generados) """

─────────────────────────────────────────────

1. CONFIGURACIÓN CELERY

─────────────────────────────────────────────

from celery import Celery, chain, group, chord from celery.utils.log import get_task_logger

app = Celery( "cultiva_ia", broker="redis://localhost:6379/0", backend="redis://localhost:6379/1", )

app.conf.update( # Límites de tiempo task_time_limit=3600, # Hard limit: 1 hora task_soft_time_limit=3000, # Soft limit: 50 minutos # Fiabilidad task_acks_late=True, # ACK tras completar, no al recibir task_reject_on_worker_lost=True, worker_prefetch_multiplier=1, # Un job a la vez por worker # Reintentos task_default_retry_delay=60, task_max_retries=3, # Serialización task_serializer="json", result_serializer="json", accept_content=["json"], # Colas especializadas task_routes={ "cultiva_ia.tasks.export_*": {"queue": "exports"}, "cultiva_ia.tasks.ai_*": {"queue": "ai_generation"}, "cultiva_ia.tasks.webhook_*": {"queue": "webhooks"}, "cultiva_ia.tasks.catalog_*": {"queue": "catalog"}, }, )

logger = get_task_logger(name)

─────────────────────────────────────────────

2. MODELO DE ESTADO DE TRABAJOS

─────────────────────────────────────────────

from uuid import uuid4 from dataclasses import dataclass, field from enum import Enum from datetime import datetime from typing import Any

class JobStatus(Enum): PENDING = "pending" RUNNING = "running" SUCCEEDED = "succeeded" FAILED = "failed"

@dataclass class Job: id: str task_type: str # "export_report" | "ai_batch" | "webhook" | "catalog" status: JobStatus created_at: datetime client_id: str # ID del cliente CULTIVA IA params: dict = field(default_factory=dict) started_at: datetime | None = None completed_at: datetime | None = None result: dict | None = None error: str | None = None attempts: int = 0

─────────────────────────────────────────────

3. REPOSITORIO DE TRABAJOS (PostgreSQL async)

─────────────────────────────────────────────

Esquema SQL esperado:

CREATE TABLE jobs (

id UUID PRIMARY KEY,

task_type TEXT NOT NULL,

status TEXT NOT NULL DEFAULT 'pending',

client_id UUID NOT NULL REFERENCES clients(id),

params JSONB,

created_at TIMESTAMPTZ NOT NULL DEFAULT now(),

started_at TIMESTAMPTZ,

completed_at TIMESTAMPTZ,

result JSONB,

error TEXT,

attempts INT NOT NULL DEFAULT 0

);

CREATE INDEX ON jobs (client_id, status);

import asyncpg # type: ignore

class JobRepository: def init(self, pool: asyncpg.Pool): self._pool = pool

async def create(self, job: Job) -> Job:
    async with self._pool.acquire() as conn:
        await conn.execute(
            """
            INSERT INTO jobs
              (id, task_type, status, client_id, params, created_at)
            VALUES ($1, $2, $3, $4, $5::jsonb, $6)
            """,
            job.id, job.task_type, job.status.value,
            job.client_id, job.params, job.created_at,
        )
    logger.info("Job created", extra={"job_id": job.id, "type": job.task_type})
    return job

async def update_status(
    self,
    job_id: str,
    status: JobStatus,
    result: dict | None = None,
    error: str | None = None,
) -> None:
    now = datetime.utcnow()
    extra_fields = {}

    if status == JobStatus.RUNNING:
        extra_fields["started_at"] = now
    elif status in (JobStatus.SUCCEEDED, JobStatus.FAILED):
        extra_fields["completed_at"] = now

    async with self._pool.acquire() as conn:
        await conn.execute(
            """
            UPDATE jobs
               SET status       = $2,
                   result       = $3::jsonb,
                   error        = $4,
                   started_at   = COALESCE($5, started_at),
                   completed_at = COALESCE($6, completed_at),
                   attempts     = attempts + 1
             WHERE id = $1
            """,
            job_id,
            status.value,
            result,
            error,
            extra_fields.get("started_at"),
            extra_fields.get("completed_at"),
        )

    logger.info(
        "Job status updated",
        extra={"job_id": job_id, "status": status.value},
    )

async def get(self, job_id: str) -> Job | None:
    async with self._pool.acquire() as conn:
        row = await conn.fetchrow(
            "SELECT * FROM jobs WHERE id = $1", job_id
        )
    if row is None:
        return None
    return Job(
        id=row["id"],
        task_type=row["task_type"],
        status=JobStatus(row["status"]),
        created_at=row["created_at"],
        client_id=row["client_id"],
        params=row["params"] or {},
        started_at=row["started_at"],
        completed_at=row["completed_at"],
        result=row["result"],
        error=row["error"],
        attempts=row["attempts"],
    )

─────────────────────────────────────────────

4. TAREAS CELERY — CASOS DE USO CULTIVA IA

─────────────────────────────────────────────

4.1 Exportación de reportes de campañas

──────────────────────────────────────────

@app.task( bind=True, name="cultiva_ia.tasks.export_campaign_report", max_retries=3, autoretry_for=(ConnectionError, TimeoutError), queue="exports", ) def export_campaign_report( self, job_id: str, client_id: str, campaign_ids: list[str], format: str = "csv", # "csv" | "pdf" date_from: str = None, date_to: str = None, ) -> dict: """ Exporta métricas de N campañas a CSV o PDF. Puede tardar minutos para clientes con millones de eventos. """ from jobs_repo import jobs_repo_sync # Versión síncrona para Celery

# Marcar como RUNNING
jobs_repo_sync.update_status(job_id, JobStatus.RUNNING)
logger.info(f"Exporting {len(campaign_ids)} campaigns for client {client_id}")

try:
    # 1. Consultar métricas (BigQuery / PostgreSQL)
    rows = analytics_db.query_campaigns(
        client_id=client_id,
        campaign_ids=campaign_ids,
        date_from=date_from,
        date_to=date_to,
    )

    # 2. Generar archivo
    if format == "csv":
        artifact_url = report_generator.to_csv(rows, job_id)
    else:
        artifact_url = report_generator.to_pdf(rows, job_id)

    # 3. Subir a S3
    public_url = storage.upload(artifact_url, prefix=f"reports/{client_id}/")

    result = {
        "url": public_url,
        "rows": len(rows),
        "format": format,
        "expires_at": (datetime.utcnow() + timedelta(days=7)).isoformat(),
    }

    jobs_repo_sync.update_status(job_id, JobStatus.SUCCEEDED, result=result)
    return result

except Exception as exc:
    if self.request.retries >= self.max_retries:
        # Falló definitivamente → DLQ
        _send_to_dlq("export_campaign_report", job_id, exc)
        jobs_repo_sync.update_status(
            job_id, JobStatus.FAILED, error=str(exc)
        )
        return {}
    raise self.retry(exc=exc, countdown=2 ** self.request.retries * 60)

4.2 Generación de contenido AI en lote

──────────────────────────────────────────

@app.task( bind=True, name="cultiva_ia.tasks.ai_batch_generate", max_retries=2, queue="ai_generation", time_limit=1800, # 30 min máx para lotes grandes soft_time_limit=1700, ) def ai_batch_generate( self, job_id: str, client_id: str, items: list[dict], # [{"id": "p123", "name": "Camiseta azul", "category": "ropa"}] template_id: str, # Plantilla de prompt configurada por el cliente model: str = "claude-sonnet-4-6", ) -> dict: """ Genera descripciones/emails/posts para N items usando LLM. Idempotente: guarda resultados parciales y reanuda si se interrumpe. """ jobs_repo_sync.update_status(job_id, JobStatus.RUNNING)

template = prompt_registry.get(template_id, client_id=client_id)
results = []
errors = []

# Reanudar desde checkpoint si existía progreso previo
checkpoint = redis_client.get(f"checkpoint:{job_id}") or {}
processed_ids = set(checkpoint.get("processed_ids", []))

for item in items:
    if item["id"] in processed_ids:
        continue  # Idempotencia: saltar ya procesados

    try:
        content = llm_client.complete(
            prompt=template.render(item),
            model=model,
            max_tokens=512,
        )
        results.append({"id": item["id"], "content": content})
        processed_ids.add(item["id"])

        # Guardar checkpoint cada 50 items
        if len(results) % 50 == 0:
            redis_client.setex(
                f"checkpoint:{job_id}",
                3600,
                {"processed_ids": list(processed_ids)},
            )

    except RateLimitError as e:
        # Backoff y retry del item
        raise self.retry(exc=e, countdown=60)
    except LLMError as e:
        errors.append({"id": item["id"], "error": str(e)})

result = {
    "generated": len(results),
    "errors": len(errors),
    "error_details": errors[:10],  # Primeros 10 errores
    "output_key": storage.save_jsonl(results, f"ai-output/{job_id}.jsonl"),
}

# Limpiar checkpoint
redis_client.delete(f"checkpoint:{job_id}")
jobs_repo_sync.update_status(job_id, JobStatus.SUCCEEDED, result=result)
return result

4.3 Webhooks a integraciones externas

──────────────────────────────────────────

@app.task( bind=True, name="cultiva_ia.tasks.webhook_dispatch", max_retries=5, queue="webhooks", ) def webhook_dispatch( self, webhook_id: str, endpoint_url: str, payload: dict, headers: dict | None = None, client_id: str = None, ) -> dict: """ Envía webhook a integraciones externas (Slack, HubSpot, Notion, Zapier). 5 reintentos con backoff exponencial. Tras agotar: Dead Letter Queue. """ import httpx

try:
    response = httpx.post(
        endpoint_url,
        json=payload,
        headers=headers or {},
        timeout=30,
    )
    response.raise_for_status()

    logger.info(
        "Webhook delivered",
        extra={
            "webhook_id": webhook_id,
            "status_code": response.status_code,
            "endpoint": endpoint_url,
        },
    )
    return {"status": "delivered", "http_status": response.status_code}

except (httpx.ConnectError, httpx.TimeoutException) as e:
    # Error transitorio → retry con backoff exponencial
    countdown = 2 ** self.request.retries * 30  # 30s, 60s, 120s, 240s, 480s
    logger.warning(
        f"Webhook transient error, retry {self.request.retries + 1}/5",
        extra={"webhook_id": webhook_id, "error": str(e)},
    )
    raise self.retry(exc=e, countdown=countdown)

except httpx.HTTPStatusError as e:
    if e.response.status_code in (400, 401, 403, 404, 422):
        # Error permanente: no tiene sentido reintentar
        _send_to_dlq("webhook_dispatch", webhook_id, e, payload=payload)
        return {"status": "permanent_failure", "http_status": e.response.status_code}

    # 5xx → reintentar
    if self.request.retries >= self.max_retries:
        _send_to_dlq("webhook_dispatch", webhook_id, e, payload=payload)
        return {"status": "exhausted", "error": str(e)}
    raise self.retry(exc=e, countdown=2 ** self.request.retries * 30)

4.4 Pipeline ETL de catálogo de producto (chain)

──────────────────────────────────────────────────

@app.task(name="cultiva_ia.tasks.catalog_extract", queue="catalog") def catalog_extract(upload_id: str, client_id: str) -> dict: """Lee CSV subido por el cliente desde S3 y parsea filas.""" raw_path = storage.get_upload_path(upload_id) rows = csv_parser.parse(raw_path) tmp_key = f"catalog-tmp/{upload_id}/raw.jsonl" storage.save_jsonl(rows, tmp_key) logger.info(f"Extracted {len(rows)} rows for upload {upload_id}") return {"upload_id": upload_id, "client_id": client_id, "rows": len(rows), "tmp_key": tmp_key}

@app.task(name="cultiva_ia.tasks.catalog_transform", queue="catalog") def catalog_transform(extract_result: dict) -> dict: """Normaliza, valida y enriquece los registros del catálogo.""" rows = storage.load_jsonl(extract_result["tmp_key"]) normalized, errors = catalog_normalizer.run(rows) out_key = f"catalog-tmp/{extract_result['upload_id']}/normalized.jsonl" storage.save_jsonl(normalized, out_key) logger.info(f"Transformed: {len(normalized)} ok, {len(errors)} errors") return {**extract_result, "normalized": len(normalized), "errors": len(errors), "out_key": out_key}

@app.task(name="cultiva_ia.tasks.catalog_load", queue="catalog") def catalog_load(transform_result: dict, job_id: str) -> dict: """Upsert de los productos normalizados en la base de datos del cliente.""" rows = storage.load_jsonl(transform_result["out_key"]) upserted = product_db.upsert_batch( client_id=transform_result["client_id"], products=rows, idempotency_key=f"upload-{transform_result['upload_id']}", ) result = { "upserted": upserted, "errors": transform_result["errors"], "upload_id": transform_result["upload_id"], } jobs_repo_sync.update_status(job_id, JobStatus.SUCCEEDED, result=result) logger.info(f"Loaded {upserted} products for job {job_id}") return result

def start_catalog_pipeline(upload_id: str, client_id: str, job_id: str): """Encadena extract → transform → load como pipeline Celery.""" pipeline = chain( catalog_extract.s(upload_id, client_id), catalog_transform.s(), catalog_load.s(job_id=job_id), ) pipeline.apply_async()

─────────────────────────────────────────────

5. DEAD LETTER QUEUE

─────────────────────────────────────────────

def _send_to_dlq(task_name: str, item_id: str, exc: Exception, **extra) -> None: """Envía tarea fallida a la cola de mensajes muertos para inspección manual.""" dlq_entry = { "task": task_name, "item_id": item_id, "error": str(exc), "error_type": type(exc).name, "failed_at": datetime.utcnow().isoformat(), **extra, } # Publicar en cola Redis dedicada redis_client.lpush("cultiva:dlq", dlq_entry) # Alertar al equipo de operaciones alert_client.send( channel="#cultiva-alerts", message=f":warning: DLQ: {task_name} id={item_id} → {type(exc).name}: {exc}", ) logger.error("Task moved to DLQ", extra=dlq_entry)

─────────────────────────────────────────────

6. ENDPOINTS FASTAPI

─────────────────────────────────────────────

from fastapi import FastAPI, HTTPException, Depends from pydantic import BaseModel

api = FastAPI(title="CULTIVA IA — Jobs API")

class ExportRequest(BaseModel): campaign_ids: list[str] format: str = "csv" date_from: str | None = None date_to: str | None = None

class AiBatchRequest(BaseModel): items: list[dict] template_id: str model: str = "claude-sonnet-4-6"

class CatalogUploadRequest(BaseModel): upload_id: str

class JobResponse(BaseModel): job_id: str status: str poll_url: str

class JobStatusResponse(BaseModel): job_id: str status: str task_type: str created_at: datetime started_at: datetime | None = None completed_at: datetime | None = None result: dict | None = None error: str | None = None is_terminal: bool

@api.post("/clients/{client_id}/exports", response_model=JobResponse) async def start_export( client_id: str, request: ExportRequest, db=Depends(get_db), ): """Inicia exportación de reportes. Retorna job_id inmediatamente.""" job_id = str(uuid4()) repo = JobRepository(db) await repo.create(Job( id=job_id, task_type="export_campaign_report", status=JobStatus.PENDING, created_at=datetime.utcnow(), client_id=client_id, params=request.model_dump(), ))

# Encolar — no espera a que termine
export_campaign_report.delay(
    job_id=job_id,
    client_id=client_id,
    **request.model_dump(),
)

return JobResponse(
    job_id=job_id,
    status="pending",
    poll_url=f"/jobs/{job_id}",
)

@api.post("/clients/{client_id}/ai-batch", response_model=JobResponse) async def start_ai_batch(client_id: str, request: AiBatchRequest, db=Depends(get_db)): """Inicia generación de contenido AI en lote.""" job_id = str(uuid4()) repo = JobRepository(db) await repo.create(Job( id=job_id, task_type="ai_batch_generate", status=JobStatus.PENDING, created_at=datetime.utcnow(), client_id=client_id, params=request.model_dump(), )) ai_batch_generate.delay( job_id=job_id, client_id=client_id, **request.model_dump(), ) return JobResponse(job_id=job_id, status="pending", poll_url=f"/jobs/{job_id}")

@api.post("/clients/{client_id}/catalog", response_model=JobResponse) async def start_catalog_import( client_id: str, request: CatalogUploadRequest, db=Depends(get_db), ): """Inicia pipeline ETL de catálogo de producto (extract→transform→load).""" job_id = str(uuid4()) repo = JobRepository(db) await repo.create(Job( id=job_id, task_type="catalog_pipeline", status=JobStatus.PENDING, created_at=datetime.utcnow(), client_id=client_id, params={"upload_id": request.upload_id}, )) start_catalog_pipeline( upload_id=request.upload_id, client_id=client_id, job_id=job_id, ) return JobResponse(job_id=job_id, status="pending", poll_url=f"/jobs/{job_id}")

@api.get("/jobs/{job_id}", response_model=JobStatusResponse) async def get_job_status(job_id: str, db=Depends(get_db)): """Polling endpoint: devuelve estado actual del job.""" repo = JobRepository(db) job = await repo.get(job_id)

if job is None:
    raise HTTPException(404, f"Job {job_id!r} not found")

return JobStatusResponse(
    job_id=job.id,
    status=job.status.value,
    task_type=job.task_type,
    created_at=job.created_at,
    started_at=job.started_at,
    completed_at=job.completed_at,
    result=job.result if job.status == JobStatus.SUCCEEDED else None,
    error=job.error if job.status == JobStatus.FAILED else None,
    is_terminal=job.status in (JobStatus.SUCCEEDED, JobStatus.FAILED),
)

─────────────────────────────────────────────

7. LANZAR WORKERS

─────────────────────────────────────────────

Arrancar workers por cola desde la CLI:

# Workers de exportaciones (I/O bound)

celery -A cultiva_ia.resultado worker -Q exports -c 8 --loglevel=info

# Workers de generación AI (CPU/API bound)

celery -A cultiva_ia.resultado worker -Q ai_generation -c 4 --loglevel=info

# Workers de webhooks (alta concurrencia)

celery -A cultiva_ia.resultado worker -Q webhooks -c 16 --loglevel=info

# Workers de catálogo (ETL)

celery -A cultiva_ia.resultado worker -Q catalog -c 4 --loglevel=info

# Monitorización con Flower

celery -A cultiva_ia.resultado flower --port=5555

─────────────────────────────────────────────

8. EJEMPLO DE USO — FLUJO COMPLETO

─────────────────────────────────────────────

# 1. Cliente solicita exportación

POST /clients/client-abc-123/exports

{ "campaign_ids": ["c1","c2","c3"], "format": "pdf", "date_from": "2026-01-01" }

→ 202 { "job_id": "f47ac10b-...", "status": "pending", "poll_url": "/jobs/f47ac10b-..." }

# 2. Cliente hace polling cada 3 segundos

GET /jobs/f47ac10b-...

→ { "status": "running", "started_at": "2026-06-16T10:01:05Z", ... }

GET /jobs/f47ac10b-...

→ {

"status": "succeeded",

"is_terminal": true,

"result": {

"url": "https://storage.cultiva.ai/reports/client-abc-123/f47ac10b.pdf",

"rows": 124850,

"format": "pdf",

"expires_at": "2026-06-23T10:01:42Z"

}

}

─────────────────────────────────────────────

// qué_hace

Proporciona patrones y ejemplos de codigo para procesar tareas de forma asincrona en Python mediante colas de trabajos y workers.

// cómo_lo_hace

Usa Celery con Redis como broker, implementando idempotencia, reintentos con backoff exponencial y gestion de estado de trabajos via base de datos.

// ejemplo_de_uso

Úsala cuando tengas tareas pesadas (envío de emails, procesado de imágenes) que no pueden bloquear la petición HTTP. Ej.: cuando el usuario sube un vídeo, encolas el transcodificado en Celery y respondes al instante con un job_id.

// plataformas

PythonCeleryRedisRQDramatiqAWS SQSGCP Tasks
Categoría
IA, ingeniería y MLOps
Tipo
Referencia
Nivel
Intermedio
Licencia
MIT
Seguridad
seguro · riesgo bajo
Versión
1.0.0

// opiniones_de_la_comunidad

Opiniones

Cargando opiniones…

// pase_cultiva_ia

Llévate todo el arsenal con el Pase

Todas las skills, prompts y automatizaciones del catálogo en un único archivo, listas para usar: un pago, acceso de por vida y las novedades que añadamos. Sin suscripción.

Pago único · IVA incluido · pago seguro con Stripe.

Acceso inmediato · si no es lo que esperabas, te devolvemos los 10 €.