Constaia
Integrations

Celery

Analyse documents in the background with Celery and the constaia SDK, with Retry-After aware retries, rate_limit, idempotency keys and batches of up to 100.

Cette page n'est pas encore traduite dans votre langue. Voici la version anglaise.

This guide moves document analysis into Celery tasks with the official constaia SDK: a task that analyses a file from your storage (by URL or as bytes), correct retries on 429 and 5xx, a rate_limit so you don't exhaust your quota, one idempotency key per task so nothing is charged twice, and a batch variant. SDK details are in the Python guide.

Requirements

  • Celery 5 with a broker (Redis or RabbitMQ) and Python ≥ 3.10.
  • A test key ck_test_… from the dashboard. In test mode no credits are consumed and the result depends on the file name.

Installation

pip install "celery[redis]" constaia

Environment variables

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

Analysis task

One client per worker process. The task receives an https:// URL of the file (for example, a presigned URL from your storage) or reads the bytes from shared storage. The task's code sets expect and checks.

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}")  # store it in your database


def mark_rejected(document_id: int, code: str | None) -> None:
    print(f"document {document_id}: rejected ({code})")  # store it in your database


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

How the pieces fit:

  • Retries on two levels. The SDK already retries 429, 408, 5xx and network errors honouring Retry-After (here max_retries=1 to keep that first level short). If it still fails, the task is rescheduled with self.retry(): waiting retry_after for RateLimitError and with exponential backoff for APIError and APIConnectionError.
  • What is not retried. InvalidRequestError (unreadable file, unsupported type, over 20 MB…) doesn't improve with retries: the document is marked. AuthenticationError and InsufficientCreditsError are not caught: the task fails and your team should see it.
  • Idempotency per task. idempotency_key=f"document-{document_id}" is stable across Celery retries. If the original request did complete, the API returns the same response (header Idempotent-Replayed: true) without charging again for 24 h. With the same key and different content the API answers 422 idempotency_key_reused, so generate the presigned URL once when enqueuing (valid long enough to cover the retries) instead of regenerating it on every attempt. More in Idempotency.
  • File name. With bytes, filename=original_name keeps the original name: in test mode it decides the response.

Enqueue from your application:

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")

The URL must be https, reachable from the internet (no private IPs), at most 20 MB and downloadable within 15 s.

Rate limit and concurrency

The API limits requests per second per key (10/s by default, see Rate limits), shared by all your workers. In Celery, rate_limit="8/s" applies per worker process, and the SDK's max_concurrency=4 limits in-flight requests per process. With several workers, split it: for example, with 4 workers use rate_limit="2/s". Any 429 that slips through is absorbed by the SDK and self.retry().

Many documents: batches

For hundreds of documents, a batch (POST /v1/batches) is more efficient: up to 100 documents per call, always asynchronous, and a single batch.completed webhook at the end. In live mode the API checks up front that you have at least 1 credit per document (otherwise 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"))

In a batch, options.export and options.metadata apply to the whole batch (combined export with one row per document, in batch["exports"]); the other options apply to each document. Batch documents don't emit analysis.completed, but they do emit analysis.review_required and analysis.failed. When they all finish you get batch.completed with the batch object. More in Bulk batches.

If you prefer individual analyses in parallel, a Celery group works too:

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 that enqueues events

The webhook is received by your web framework, verified with the raw body and enqueued without processing. For example, with Flask (in Django and FastAPI it's the same with request.body and 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

Deliveries are retried for about 3 days with the same webhook-id: make the tasks idempotent or deduplicate on that id. Format and retries in Webhooks.

Errors

HTTPExceptionIn the task
400, 413, 415, 422InvalidRequestErrorDon't retry; mark the document (exc.code)
401AuthenticationErrorLet it fail: configuration error
402InsufficientCreditsErrorLet it fail and alert; retry once credits are available
429RateLimitErrorself.retry(countdown=exc.retry_after)
5xxAPIErrorself.retry() with exponential backoff
—APIConnectionError, APITimeoutErrorself.retry() with exponential backoff

Always log exc.request_id. Codes in Errors.

Tests

With task_always_eager tasks run in the same process. With CONSTAIA_API_KEY=ck_test_... the API answers based on the file name and charges nothing; copy any real JPEG into tests/fixtures/ as dni_valid.jpg and 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()

More test files in Test mode.

Production checklist

  • soft_time_limit above 60 s (the SDK waits up to 60 s per attempt) and the broker's visibility_timeout larger than the task's maximum time including retries.
  • rate_limit × workers and max_concurrency × processes below your key's limit.
  • Stable idempotency key per document or batch; presigned URL generated once at enqueue time.
  • Don't retry InvalidRequestError; alert on AuthenticationError and InsufficientCreditsError.
  • The route that receives user files requires authentication and rate limiting, even though the analysis runs in the background: every live analysis consumes credits.
  • Webhook verified with the raw body that only enqueues; idempotent tasks.
  • CONSTAIA_API_KEY=ck_live_... only in production. Read Storage and privacy: with batches the file is stored encrypted only until the analysis finishes when storage: none.

Next steps

Sur cette page