Airflow and Python pipelines
Two ways to decide over many records in a pipeline:
curva map |
AsyncCurva.decide_many |
|
|---|---|---|
| Runs | the CLI, straight against the model provider | the Python SDK, against a Curva server |
| Input | a JSONL file | any Python iterable |
| Resumable | yes, reruns continue after the last line written | no, your task retries decide |
| Calibration, audit log, API keys | no (no server involved) | yes |
Use curva map for one-off or offline backfills of files, and decide_many when the answers
should be calibrated by feedback and appear in the audit log and the dashboard.
curva map in a task
Section titled “curva map in a task”export OPENROUTER_API_KEY=sk-or-v1-...curva map tickets.jsonl -q questions.json -o answers.jsonl --model model-a,model-bBecause a rerun continues after the last line written, a retried Airflow task picks up where the failed one stopped without paying twice. In Airflow:
from airflow.operators.bash import BashOperator
classify = BashOperator( task_id="classify_tickets", bash_command="curva map /data/{{ ds }}/tickets.jsonl -q /opt/curva/questions.json -o /data/{{ ds }}/answers.jsonl", env={"OPENROUTER_API_KEY": "{{ var.value.openrouter_key }}"}, append_env=True, retries=3,)Input and output formats are in Batch files with curva map.
decide_many against a server
Section titled “decide_many against a server”import asynciofrom curva import AsyncCurva, Choice, Noul
QUESTIONS = { "department": Choice("Which team should handle this", ["billing", "technical", "sales"], min_confidence=0.8), "refund": Noul("The customer explicitly asks for a refund"),}
async def classify(rows): async with AsyncCurva() as curva: # CURVA_BASE_URL, CURVA_API_KEY results = await curva.decide_many( [({"ticket": r["text"]}, QUESTIONS) for r in rows], concurrency=16, project="support", return_exceptions=True, ) out = [] for row, d in zip(rows, results): if isinstance(d, Exception): out.append({**row, "error": str(d)}) continue team = d["department"] out.append({**row, "decision_id": d.id, "team": team.choice, "confidence": team.confidence, "needs_human": bool(team.abstain), "refund": d["refund"].noul}) return outResults keep the input order. return_exceptions=True puts each failure in its slot, so one bad
record doesn’t fail the batch. The client already retries 429 and 5xx (honouring
Retry-After), and the server’s per-key rate limit keeps a large batch polite.
In Airflow, wrap it in a task:
from airflow.decorators import task
@task(retries=2)def classify_tickets(rows: list[dict]) -> list[dict]: return asyncio.run(classify(rows))Store decision_id with each row. When labels arrive (a later DAG, a review table), send them
with Curva().feedback(decision_id, "department", label); see
Calibrate with feedback.
Other orchestrators
Section titled “Other orchestrators”The same two patterns work in Dagster, Prefect, Luigi or a cron job: shell out to curva map
for files, or call decide_many from Python. The SDK has no dependencies, so it installs cleanly
into any worker image (pip install curva-ai).

