Playbook · 6 minute read
How to Build an Airflow AI Pipeline
An Airflow AI pipeline schedules batch inference, embedding refresh, evaluation runs, and backfills as DAGs whose tasks push model calls to executors or external services rather than the scheduler, are idempotent so retries and backfills are safe, record model and prompt versions per run, and schedule expensive work into cost-appropriate windows.
The AI work that gets the attention is interactive: assistants, agents, and real-time features. The AI work that consumes most of the volume is batch: classifying last night's documents, refreshing embeddings after content changed, running evaluation sets against a new prompt, backfilling a million records after a model upgrade. Airflow already orchestrates the data platform's batches, and it runs these well, provided model calls sit where they belong and tasks are built to be rerun. This guide covers how, drawing on FISTA Solutions' AI enablement delivery on data platforms. It complements airflow vs dagster and batch vs real-time inference.
Which AI workloads belong in Airflow?
| Workload | Shape | Notes |
|---|---|---|
| Batch classification and extraction | Nightly or hourly over new records | Highest volume, clearest return |
| Embedding generation and refresh | Triggered by content change or scheduled | Keeps retrieval indexes current |
| Evaluation runs | On change and on schedule | Regression detection |
| Backfills | One-off after model or prompt change | Must be idempotent |
| Data preparation for tuning | Periodic | Feeds training pipelines |
| Cost and usage reporting | Daily | Aggregates gateway logs |
| Index rebuilds | Scheduled or triggered | Retrieval infrastructure maintenance |
Interactive workloads do not belong here; Airflow is a scheduler, not a request handler.
Where should model calls run?
Not in the scheduler, and this is the mistake that causes the most damage. Airflow's scheduler must stay responsive to schedule every DAG; a task that makes thousands of slow model calls in scheduler-side code stalls it.
Model calls run in executor tasks, with concurrency bounded by pools sized to provider rate limits, or are delegated to an external batch inference service that a task submits work to and another task polls for completion. The second pattern suits large volumes: the DAG orchestrates, the batch service does the inference, and the scheduler is never blocked.
How are tasks made idempotent?
By keying outputs on input identity plus model and prompt version, checking for existing results before inferring, and writing results in an overwrite-safe way. A retried task, a backfill over a historical range, or an operator's manual rerun then produces the same state rather than duplicated records or repeated cost.
Idempotency also enables partial reprocessing: when a prompt changes, only records without results for the new version need inference, which turns a full backfill into an incremental one.
How should evaluation run?
As a DAG. The evaluation DAG loads the reference set, runs it against the current model and prompt through the same path production uses, scores results, compares against the previous run, and fails or alerts on regression beyond a threshold. It triggers on every change to prompt, model, or retrieval configuration, and runs on a schedule to catch drift from data changes nobody flagged.
Making evaluation a first-class DAG means it appears in the same monitoring, has the same retry semantics, and produces results in tables that dashboards read, alongside the pipelines it protects. See the RAG evaluation methodology whitepaper.
What should each run record?
Model identifier and version, prompt version, retrieval index version where relevant, the input data snapshot or range, the run's token and cost totals, and the output location. Without those, a quality change six weeks later cannot be attributed to any cause, and a result cannot be reproduced.
These belong in the run's metadata and in the output records themselves, so downstream consumers of a classification know which model produced it.
How is cost controlled?
Through several mechanisms that compound. Batch pricing, which several providers offer at a substantial discount for asynchronous workloads with relaxed latency, suits most Airflow AI work. Off-peak scheduling avoids contention. Pool-based concurrency limits prevent a backfill from consuming a provider quota that interactive systems need. Routing bounded tasks to smaller models cuts cost per record. And recording tokens and cost per task makes every DAG's spend visible next to its runtime.
The DAG that classifies a million records nightly is usually the largest single AI cost line in an organisation, and it is the easiest to optimise once it is measured. See the AI FinOps whitepaper.
How are embeddings kept current?
Through a refresh DAG triggered by content change events or scheduled at a cadence matched to content velocity, which embeds only changed or new content, writes to the index with versioning, and handles deletions so removed content leaves the index. Full re-embedding is reserved for embedding model changes and is a backfill with the same idempotency discipline.
How are failures handled?
With bounded retries tuned to model failure modes, backoff on rate limits, and a failure path that records which inputs failed so they can be reprocessed rather than lost. Task-level failure should not fail the whole batch: a DAG that processes ten thousand records should complete for the nine thousand that succeeded and report the rest.
What does the build sequence look like?
One week establishing the pattern: a classification DAG with executor-side inference, pool limits, idempotent outputs, and version recording. One week on the evaluation DAG triggered by change. One week on embedding refresh. Then backfill tooling and cost reporting, which reuse the pattern.
What goes wrong?
Model calls in the scheduler. Non-idempotent tasks, discovered on the first backfill. Evaluation run by hand or not at all. Versions not recorded. Concurrency unbounded, exhausting provider quotas. Full re-embedding on every content change. And batch pricing never enabled because nobody checked whether the workload tolerated the latency, which nearly all batch work does.
How does this fit with a dedicated ML orchestrator?
Often alongside rather than instead. Organisations with a model training platform may keep training and fine-tuning pipelines there, while inference batches, embedding refresh, and evaluation runs live in Airflow next to the data pipelines they depend on. The dividing line that works is data proximity: work that reads and writes the same tables as the rest of the platform belongs in the platform's scheduler, and work that is primarily about model artifacts belongs with the model tooling. Running inference batches in a separate orchestrator from the data pipelines that feed them produces the scheduling gaps and dependency mismatches that cause most batch AI incidents.
How FISTA Solutions helps
FISTA Solutions builds Airflow AI pipelines with inference in executor tasks or external batch services, idempotent outputs keyed by version, evaluation as a first-class DAG, incremental embedding refresh, and cost recorded per task, through AI enablement, AI agents, and forward deployed engineers working with data platform teams. The record behind the approach is 150+ projects for 50+ companies with 99.9% uptime.
To run batch AI at scale without stalling the scheduler, message FISTA on WhatsApp, or read batch vs real-time inference.
Share-ready article cover
Download the generated social format.
Clear answers
Questions raised by this field note.
Straightforward guidance for evaluating scope, fit, and the next step.
01What AI workloads suit Airflow?
Batch inference over document and record sets, embedding generation and refresh for retrieval indexes, scheduled evaluation runs against reference sets, backfills when a model or prompt changes, and data preparation for training or fine-tuning, all of which are scheduled, bounded, and naturally expressed as DAGs.
02Where should model calls run?
In executor tasks with bounded concurrency, or delegated to an external batch inference service that the task submits to and polls, never in scheduler-side code. Long-running model calls in the wrong place stall the scheduler and starve every other DAG.
03How is idempotency achieved?
By keying inference outputs on input identity and model version, checking for existing results before calling the model, and writing results in a way that overwrites rather than duplicates, so a retried task, a backfill, or a manual rerun produces the same state.
04How does evaluation fit?
As a DAG that runs reference sets against the current model and prompt, scores results, compares with the previous run, and fails or alerts on regression, triggered on every prompt or model change and on a schedule to catch drift from data changes.
05How is cost controlled?
Through batch pricing where providers offer it, scheduling expensive runs into off-peak windows, concurrency limits per pool, routing to cheaper model tiers for bounded tasks, and recording tokens and cost per task so each DAG's cost is visible alongside its runtime.
Continue exploring
Related capabilities
Start with the hard problem
Need the outcome owned, not merely analyzed?
Tell us where delivery is constrained. We’ll map the fastest credible path from intent to verified production.