Constaia
Integraciones

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]" constaia

Variables de entorno

.env
CONSTAIA_API_KEY=ck_test_...
CONSTAIA_WEBHOOK_SECRET=whsec_...
CELERY_BROKER_URL=redis://localhost:6379/0

Tarea 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.

tasks.py
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, 5xx y errores de red respetando Retry-After (aquí max_retries=1 para que el primer nivel sea corto). Si aun así falla, la tarea se reprograma con self.retry(): con la espera de retry_after para RateLimitError y con backoff exponencial para APIError y APIConnectionError.
  • Lo que no se reintenta. InvalidRequestError (fichero ilegible, tipo no admitido, más de 20 MB…) no mejora reintentando: se marca el documento. AuthenticationError e InsufficientCreditsError no 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 (cabecera Idempotent-Replayed: true) sin volver a cobrar durante 24 h. Con la misma clave y un contenido distinto la API responde 422 idempotency_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_name conserva el nombre original: en modo test decide la respuesta.

Encola desde tu aplicación:

enqueue.py
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).

batch_tasks.py
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()):

webhook_app.py
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 "", 204

Las 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

HTTPExcepciónEn la tarea
400, 413, 415, 422InvalidRequestErrorNo reintentar; marcar el documento (exc.code)
401AuthenticationErrorDejar fallar: error de configuración
402InsufficientCreditsErrorDejar fallar y avisar; reintentar cuando haya créditos
429RateLimitErrorself.retry(countdown=exc.retry_after)
5xxAPIErrorself.retry() con backoff exponencial
—APIConnectionError, APITimeoutErrorself.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.

tests/test_tasks.py
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_limit por encima de 60 s (el SDK espera hasta 60 s por intento) y visibility_timeout del broker mayor que el tiempo máximo de la tarea con sus reintentos.
  • rate_limit × workers y max_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 en AuthenticationError e InsufficientCreditsError.
  • 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 con storage: none.

Siguientes pasos

En esta página