Celery
Analiza documentos en segundo plano con Celery y el SDK constaia, con reintentos según Retry-After, rate_limit, idempotencia y lotes de hasta 100.
Esta guía mueve el análisis de documentos a tareas de Celery con el SDK oficial constaia: una tarea que analiza un
fichero desde tu almacenamiento (por URL o en bytes), reintentos correctos ante 429 y 5xx, un rate_limit para no
saturar tu cuota, una clave de idempotencia por tarea para no cobrar dos veces y una variante con
lotes. Los detalles del SDK están en la guía de Python.
Requisitos
- Celery 5 con un broker (Redis o RabbitMQ) y Python ≥ 3.10.
- Una clave de test
ck_test_…del panel. En modo test no se consumen créditos y el resultado depende del nombre del fichero.
Instalación
pip install "celery[redis]" constaiaVariables de entorno
CONSTAIA_API_KEY=ck_test_...
CONSTAIA_WEBHOOK_SECRET=whsec_...
CELERY_BROKER_URL=redis://localhost:6379/0Tarea de análisis
Un cliente por proceso worker. La tarea recibe una URL https:// del fichero (por ejemplo, una URL prefirmada de tu
almacenamiento) o lee los bytes de un almacenamiento compartido. expect y checks los fija el código de la tarea.
import os
from pathlib import Path
from celery import Celery
from constaia import (
APIConnectionError,
APIError,
Constaia,
InvalidRequestError,
RateLimitError,
)
app = Celery("documents", broker=os.environ["CELERY_BROKER_URL"])
app.conf.task_acks_late = True
client = Constaia(timeout=60.0, max_retries=1, max_concurrency=4)
IDENTITY = {
"expect": ["es_dni", "es_nie", "passport"],
"checks": {"not_expired": True, "min_age_years": 18},
"storage": "none",
}
def save_result(document_id: int, analysis: dict) -> None:
status = (analysis.get("verdict") or {}).get("status") or analysis["status"]
print(f"document {document_id}: {analysis['id']} {status}") # guarda en tu base de datos
def mark_rejected(document_id: int, code: str | None) -> None:
print(f"document {document_id}: rejected ({code})") # guarda en tu base de datos
@app.task(bind=True, rate_limit="8/s", max_retries=5, soft_time_limit=180)
def analyze_document_url(self, document_id: int, file_url: str) -> str:
try:
analysis = client.analyze(
file_url=file_url,
**IDENTITY,
metadata={"document_id": str(document_id)},
idempotency_key=f"document-{document_id}",
)
except RateLimitError as exc:
raise self.retry(exc=exc, countdown=exc.retry_after or 5)
except (APIError, APIConnectionError) as exc:
raise self.retry(exc=exc, countdown=min(30 * 2**self.request.retries, 900))
except InvalidRequestError as exc:
mark_rejected(document_id, exc.code)
return "rejected"
save_result(document_id, analysis)
return analysis["id"]
@app.task(bind=True, rate_limit="8/s", max_retries=5, soft_time_limit=180)
def analyze_document_file(self, document_id: int, path: str, original_name: str) -> str:
try:
analysis = client.analyze(
Path(path).read_bytes(),
filename=original_name,
**IDENTITY,
metadata={"document_id": str(document_id)},
idempotency_key=f"document-{document_id}",
)
except RateLimitError as exc:
raise self.retry(exc=exc, countdown=exc.retry_after or 5)
except (APIError, APIConnectionError) as exc:
raise self.retry(exc=exc, countdown=min(30 * 2**self.request.retries, 900))
except InvalidRequestError as exc:
mark_rejected(document_id, exc.code)
return "rejected"
save_result(document_id, analysis)
Path(path).unlink(missing_ok=True)
return analysis["id"]Cómo encaja cada pieza:
- Reintentos en dos niveles. El SDK ya reintenta
429,408,5xxy errores de red respetandoRetry-After(aquímax_retries=1para que el primer nivel sea corto). Si aun así falla, la tarea se reprograma conself.retry(): con la espera deretry_afterparaRateLimitErrory con backoff exponencial paraAPIErroryAPIConnectionError. - Lo que no se reintenta.
InvalidRequestError(fichero ilegible, tipo no admitido, más de 20 MB…) no mejora reintentando: se marca el documento.AuthenticationErroreInsufficientCreditsErrorno se capturan: la tarea falla y debe verlo tu equipo. - Idempotencia por tarea.
idempotency_key=f"document-{document_id}"es estable entre reintentos de Celery. Si la petición original llegó a completarse, la API devuelve la misma respuesta (cabeceraIdempotent-Replayed: true) sin volver a cobrar durante 24 h. Con la misma clave y un contenido distinto la API responde422idempotency_key_reused, así que genera la URL prefirmada una vez al encolar (con validez suficiente para los reintentos) en lugar de regenerarla en cada intento. Más en Idempotencia. - Nombre del fichero. Con bytes,
filename=original_nameconserva el nombre original: en modo test decide la respuesta.
Encola desde tu aplicación:
from tasks import analyze_document_file, analyze_document_url
analyze_document_url.delay(42, "https://storage.example.com/uploads/dni_valid.jpg?X-Signature=...")
analyze_document_file.delay(43, "/srv/shared/uploads/7f3a.jpg", "dni_valid.jpg")La URL debe ser https, accesible desde internet (sin IPs privadas), de 20 MB como máximo y descargable en 15 s.
Rate limit y concurrencia
La API limita las peticiones por segundo por clave (10/s por defecto, ver
Rate limits), compartidas por todos tus workers. En Celery, rate_limit="8/s" se
aplica por proceso worker, y max_concurrency=4 del SDK limita las peticiones en vuelo por proceso. Con varios
workers, reparte: por ejemplo, con 4 workers usa rate_limit="2/s". Los 429 que se escapen los absorben el SDK y
self.retry().
Muchos documentos: lotes
Para cientos de documentos, un lote (POST /v1/batches) es más eficiente: hasta 100 documentos por llamada, siempre
asíncrono, y un solo webhook batch.completed al terminar. En modo live la API comprueba de antemano que tengas al
menos 1 crédito por documento (si no, 402).
from constaia import APIConnectionError, APIError, RateLimitError
from tasks import app, client
CHUNK = 100
@app.task(bind=True, max_retries=5)
def create_batch(self, job_id: str, chunk_index: int, file_urls: list[str]) -> str:
try:
batch = client.batches.create(
items=[{"file_url": url} for url in file_urls],
options={"expect": "invoice", "export": ["xlsx"], "metadata": {"job_id": job_id}},
idempotency_key=f"batch-{job_id}-{chunk_index}",
)
except RateLimitError as exc:
raise self.retry(exc=exc, countdown=exc.retry_after or 5)
except (APIError, APIConnectionError) as exc:
raise self.retry(exc=exc, countdown=min(30 * 2**self.request.retries, 900))
return batch["id"]
def submit_job(job_id: str, file_urls: list[str]) -> None:
for index in range(0, len(file_urls), CHUNK):
create_batch.delay(job_id, index // CHUNK, file_urls[index:index + CHUNK])
@app.task
def handle_batch_completed(batch: dict) -> None:
print(batch["id"], batch["counts"])
for analysis_id in batch["analyses"]:
analysis = client.analyses.get(analysis_id)
status = (analysis.get("verdict") or {}).get("status") or analysis["status"]
print(analysis_id, status, analysis["fields"].get("total", {}).get("value"))En un lote, options.export y options.metadata se aplican al lote entero (exportación combinada con una fila por
documento, en batch["exports"]); el resto de opciones se aplica a cada documento. Los documentos de un lote no
emiten analysis.completed, pero sí analysis.review_required y analysis.failed. Cuando todos terminan llega
batch.completed con el objeto lote. Más en Lotes masivos.
Si prefieres análisis individuales en paralelo, un group de Celery también sirve:
from celery import group
from tasks import analyze_document_url
group(analyze_document_url.s(doc_id, url) for doc_id, url in pending).apply_async()Webhook que encola eventos
El webhook se recibe en tu framework web, se verifica con el cuerpo crudo y se encola sin procesar. Por ejemplo, con
Flask (en Django y FastAPI es igual con request.body y
await request.body()):
import os
from flask import Flask, abort, request
import constaia
from batch_tasks import handle_batch_completed
from tasks import app as celery_app
web = Flask(__name__)
@celery_app.task
def handle_analysis_event(event: dict) -> None:
analysis = event["data"]
print(event["type"], analysis["id"], (analysis.get("verdict") or {}).get("status"))
@web.post("/webhooks/constaia")
def constaia_webhook():
try:
event = constaia.webhooks.verify(request.get_data(), request.headers, os.environ["CONSTAIA_WEBHOOK_SECRET"])
except constaia.WebhookVerificationError:
abort(400)
if event["type"] == "batch.completed":
handle_batch_completed.delay(event["data"])
elif event["type"] in ("analysis.completed", "analysis.review_required", "analysis.failed"):
handle_analysis_event.delay(event)
return "", 204Las entregas se reintentan hasta unos 3 días con el mismo webhook-id: haz las tareas idempotentes o deduplica por
ese id. Formato y reintentos en Webhooks.
Errores
| HTTP | Excepción | En la tarea |
|---|---|---|
| 400, 413, 415, 422 | InvalidRequestError | No reintentar; marcar el documento (exc.code) |
| 401 | AuthenticationError | Dejar fallar: error de configuración |
| 402 | InsufficientCreditsError | Dejar fallar y avisar; reintentar cuando haya créditos |
| 429 | RateLimitError | self.retry(countdown=exc.retry_after) |
| 5xx | APIError | self.retry() con backoff exponencial |
| — | APIConnectionError, APITimeoutError | self.retry() con backoff exponencial |
Registra siempre exc.request_id. Códigos en Errores.
Tests
Con task_always_eager las tareas se ejecutan en el mismo proceso. Con CONSTAIA_API_KEY=ck_test_... la API responde
según el nombre del fichero y no cobra; copia cualquier JPEG real a tests/fixtures/ como dni_valid.jpg y
dni_expired.jpg.
import os
import shutil
from pathlib import Path
import pytest
from tasks import analyze_document_file, app
FIXTURES = Path(__file__).parent / "fixtures"
pytestmark = pytest.mark.skipif(
not os.environ.get("CONSTAIA_API_KEY", "").startswith("ck_test_"), reason="needs a ck_test_ key"
)
@pytest.fixture(autouse=True)
def eager():
app.conf.task_always_eager = True
yield
app.conf.task_always_eager = False
@pytest.mark.parametrize("name", ["dni_valid.jpg", "dni_expired.jpg"])
def test_analyze_document_file(tmp_path, name):
stored = tmp_path / "upload.bin"
shutil.copy(FIXTURES / name, stored)
result = analyze_document_file.delay(1, str(stored), name).get()
assert result.startswith("an_")
assert not stored.exists()Más ficheros de prueba en Modo test.
Checklist de producción
soft_time_limitpor encima de 60 s (el SDK espera hasta 60 s por intento) yvisibility_timeoutdel broker mayor que el tiempo máximo de la tarea con sus reintentos.rate_limit× workers ymax_concurrency× procesos por debajo del límite de tu clave.- Clave de idempotencia estable por documento o lote; URL prefirmada generada una vez al encolar.
- No reintentar
InvalidRequestError; alertas enAuthenticationErroreInsufficientCreditsError. - La ruta que recibe los ficheros de usuarios exige autenticación y límite de frecuencia, aunque el análisis sea en segundo plano: cada análisis live consume créditos.
- Webhook verificado con el cuerpo crudo que solo encola; tareas idempotentes.
CONSTAIA_API_KEY=ck_live_...solo en producción. Revisa Almacenamiento y privacidad: con lotes el fichero se guarda cifrado solo hasta que termina el análisis constorage: none.
Siguientes pasos
FastAPI
Valida documentos en FastAPI con el cliente async de constaia, UploadFile, modelos pydantic de respuesta, BackgroundTasks y un webhook que verifica la firma.
Ruby on Rails
Valida DNI y otros documentos desde Rails 7/8 con Faraday multipart, un service object, ActiveJob y webhooks de Constaia verificados con HMAC.