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.
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]" constaiaEnvironment variables
CONSTAIA_API_KEY=ck_test_...
CONSTAIA_WEBHOOK_SECRET=whsec_...
CELERY_BROKER_URL=redis://localhost:6379/0Analysis 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.
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,5xxand network errors honouringRetry-After(heremax_retries=1to keep that first level short). If it still fails, the task is rescheduled withself.retry(): waitingretry_afterforRateLimitErrorand with exponential backoff forAPIErrorandAPIConnectionError. - What is not retried.
InvalidRequestError(unreadable file, unsupported type, over 20 MB…) doesn't improve with retries: the document is marked.AuthenticationErrorandInsufficientCreditsErrorare 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 (headerIdempotent-Replayed: true) without charging again for 24 h. With the same key and different content the API answers422idempotency_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_namekeeps the original name: in test mode it decides the response.
Enqueue from your application:
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).
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()):
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 "", 204Deliveries 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
| HTTP | Exception | In the task |
|---|---|---|
| 400, 413, 415, 422 | InvalidRequestError | Don't retry; mark the document (exc.code) |
| 401 | AuthenticationError | Let it fail: configuration error |
| 402 | InsufficientCreditsError | Let it fail and alert; retry once credits are available |
| 429 | RateLimitError | self.retry(countdown=exc.retry_after) |
| 5xx | APIError | self.retry() with exponential backoff |
| — | APIConnectionError, APITimeoutError | self.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.
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_limitabove 60 s (the SDK waits up to 60 s per attempt) and the broker'svisibility_timeoutlarger than the task's maximum time including retries.rate_limit× workers andmax_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 onAuthenticationErrorandInsufficientCreditsError. - 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 whenstorage: none.
Next steps
FastAPI
Validate documents in FastAPI with the async constaia client, UploadFile, pydantic response models, BackgroundTasks and a webhook that verifies the signature.
Ruby on Rails
Validate Spanish IDs and other documents from Rails 7/8 with Faraday multipart, a service object, ActiveJob and HMAC-verified Constaia webhooks.