Repository navigation
A hard dollar cap per DAG run for mapped tasks that call OpenAI #74174
Unanswered
domondi1
asked this question in
Show and tell
Replies: 0 comments
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Uh oh!
There was an error while loading. Please reload this page.
A common LLM pattern in Airflow is
summarize.expand(item=...): one DAG run maps a task over hundreds of rows, each mapped task makes a model call, and several run at once. Nothing in that setup stops a single run at a dollar amount. A bad upstream extract (10x the rows) or a prompt change that blows up output length shows up on the provider bill, not in the run.Here's a way to cap one DAG run's spend using only the OpenAI provider as it is. Point the
openaiconnection at an OpenAI-compatible gateway that enforces per-run budgets, and send the DAG run's id and budget as headers on each call. Every mapped task instance in that run then draws from one budget, even across workers.Run the gateway somewhere your workers can reach (it uses your provider key):
Connection (env var form;
hostbecomes the client'sbase_url):After the run:
shows how many calls the run made and what they cost.
Each call reserves its worst-case cost before it's sent, so concurrent mapped tasks can't overspend together. Once the run's budget can't cover another call, the call gets HTTP 402 before it reaches OpenAI. The task fails with
openai.APIStatusError(status 402) and the rest of the run is protected. If you setretrieson the task, each retry is refused the same way without spending anything, but you'll probably wantretries=0on these tasks or a check for 402. The headers aren't forwarded to OpenAI.What I tested: Airflow 3.3.2 with apache-airflow-providers-openai 2.0.0 and inferrail 0.4.12, using
airflow dags testand a local stub upstream standing in for OpenAI, with a 10-way mapped task. With a generous budget, all 10 succeeded andinferrail work <run_id>showed 10 calls. With a tight one, 3 succeeded, 7 failed with 402 before reaching the upstream, and the run's total stayed under the cap. I haven't run it on a deployed Airflow against the real API yet, so reports are welcome.Caveats: chat completions only; the budget store is a SQLite file on the gateway host, so run one gateway for all workers; and the model needs a price in Inferrail (
gpt-4o-mini,gpt-4.1-mini,gpt-4.1built in,inferrail modelslists others, and you can add your own prices). To budget per mapped item instead of per run, usef"{ctx['run_id']}:{ctx['ti'].map_index}"as the id.Setup details: guide. I maintain Inferrail (open source, Apache-2.0), so take the suggestion with that in mind.
All reactions