Skip to content

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.

Terminal window
export OPENROUTER_API_KEY=sk-or-v1-...
curva map tickets.jsonl -q questions.json -o answers.jsonl --model model-a,model-b

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

import asyncio
from 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 out

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

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

© 2026 Tarkova Private Limited.