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.
Descarga abierta · sin registro · para Python, Celery, Redis
""" 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
// 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 €.