GitHub user domondi1 created a discussion: A hard dollar cap per DAG run for mapped tasks that call OpenAI
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 `openai` connection 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): ```bash pip install inferrail export OPENAI_API_KEY=sk-... inferrail serve --quickstart --app-mode # http://127.0.0.1:8000 ``` Connection (env var form; `host` becomes the client's `base_url`): ```bash export AIRFLOW_CONN_INFERRAIL='{"conn_type": "openai", "password": "unused", "host": "http://127.0.0.1:8000/v1"}' ``` ```python import pendulum from airflow.sdk import dag, task, get_current_context from airflow.providers.openai.hooks.openai import OpenAIHook @dag(schedule=None, start_date=pendulum.datetime(2026, 1, 1), catchup=False, params={"budget_usd": "5.00"}) def ticket_summaries(): @task def tickets(): return load_ticket_texts() @task(max_active_tis_per_dagrun=10) def summarize(ticket: str): ctx = get_current_context() client = OpenAIHook(conn_id="inferrail").get_conn() r = client.chat.completions.create( model="gpt-4o-mini", max_tokens=800, # the reservation is based on it extra_headers={ "X-Inferrail-Attribute-Work-Id": ctx["run_id"], "X-Inferrail-Budget-Usd": ctx["params"]["budget_usd"], }, messages=[{"role": "user", "content": f"Summarize: {ticket}"}], ) return r.choices[0].message.content summarize.expand(ticket=tickets()) ticket_summaries() ``` After the run: ```bash inferrail work "manual__2026-10-03T20:52:40.229888+00:00" # the DAG run id ``` 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 set `retries` on the task, each retry is refused the same way without spending anything, but you'll probably want `retries=0` on 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 test` and a local stub upstream standing in for OpenAI, with a 10-way mapped task. With a generous budget, all 10 succeeded and `inferrail 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.1` built in, `inferrail models` lists others, and you can add your own prices). To budget per mapped item instead of per run, use `f"{ctx['run_id']}:{ctx['ti'].map_index}"` as the id. Setup details: [guide](https://tryinferrail.com/recipes/agent-run-budget.html?ref=airflow-discussions). I maintain Inferrail (open source, Apache-2.0), so take the suggestion with that in mind. GitHub link: https://github.com/apache/airflow/discussions/74174 ---- This is an automatically sent email for [email protected]. To unsubscribe, please send an email to: [email protected]
