Sam Austin AI

Apache Airflow Tutorial: Orchestrate Your ML Workflows (2026)

September 8, 2026 15 min read Sam Austin
Contents
Apache Airflow ML Pipeline Orchestration
Apache Airflow ML Pipeline Orchestration

Figure 1: Apache Airflow's web UI provides real-time monitoring of ML pipeline execution, task status, and failure diagnostics

The ETL pipeline article's honest verdict on when plain Python stops cutting it applies directly to ML workflows: multi-step dependencies, failure recovery, scheduled reruns. ML pipelines hit that wall faster than almost any other workload — a training run that depends on fresh data, which depends on a successful extraction, which might fail at 3am for reasons that have nothing to do with your model. Airflow exists specifically for orchestrating exactly that kind of interdependent mess reliably.

Currently at version 3.3.1, and now explicitly describing itself as covering "traditional time-based or event-triggered batch-oriented data pipelines, machine learning, model training, and agentic or LLM-based workloads" — Airflow's own positioning has genuinely expanded to claim the ML and agentic space directly, not just classic ETL. Worth knowing if you last touched this tool years ago and think of it as purely a data-engineering thing.

By the end of this guide, you'll understand DAGs, tasks, and operators well enough to build a real ML retraining pipeline, and you'll know the specific best practices that separate a DAG that works in testing from one that survives production. IMO, this is genuinely the piece of infrastructure that turns "I have a Jupyter notebook that retrains my model" into "I have a system that retrains my model" :)

The Five Concepts You Actually Need First

Before writing any code, these five ideas are the building blocks of literally every Airflow pipeline, ML or otherwise.

DAG (Directed Acyclic Graph) — a Python file defining your workflow: which tasks exist, what order they run in, and what depends on what. "Directed" means tasks flow one way; "acyclic" means no loops — Task A can't depend on Task B if Task B depends on Task A.

Task — a single unit of work inside a DAG: running a SQL query, calling an API, sending an alert, training a model.

Operator — defines what a task actually does; Airflow ships dozens of built-in operators, and tasks are simply instances of them.

Schedule — when the workflow runs: a cron expression, a fixed interval, or (increasingly, as of recent versions) an event trigger like a new file landing in a bucket.

XComs — Airflow's mechanism for passing small pieces of data between tasks, like a run ID or a model's evaluation score, without shoving your entire dataset through the orchestrator itself.

Installing Airflow

python -m venv airflow-env
source airflow-env/bin/activate

pip install apache-airflow

airflow db migrate
airflow standalone

airflow standalone genuinely spins up the scheduler, webserver, and a default admin account all at once — the fastest path to a working local instance for learning purposes. The web UI lands at localhost:8080, where you can trigger DAGs manually, inspect logs, and monitor task status without touching the command line again.

Why ML Workflows Specifically Suit Airflow

This isn't a niche use case being retrofitted — community surveys have found ML and data science pipelines to be the second most common Airflow use case after general ETL/data ingestion, and the tool's Python-native DAG definition makes it a genuinely natural fit for data scientists and ML engineers who already think in Python rather than a separate DSL.

A typical ML pipeline DAG breaks down into a recognizable sequence: fetch new data, validate it, engineer features, train the model, evaluate it against a threshold, and — only if it passes — register and deploy it. Each of those is a task; the arrows between them are dependencies Airflow enforces automatically.

Older Airflow tutorials show you PythonOperator instances wired together explicitly. The modern, recommended approach uses the @task decorator instead — genuinely cleaner and less repetitive.

from airflow.decorators import dag, task
from datetime import datetime
import pandas as pd

@dag(schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False)
def ml_retraining_pipeline():

    @task
    def fetch_new_data():
        df = pd.read_csv("https://example.com/latest_sales.csv")
        path = "/tmp/raw_data.parquet"
        df.to_parquet(path)
        return path

    @task
    def validate_data(data_path: str):
        df = pd.read_parquet(data_path)
        if df.isnull().sum().sum() > len(df) * 0.1:
            raise ValueError("Too many missing values — halting pipeline")
        return data_path

    @task
    def train_model(data_path: str):
        df = pd.read_parquet(data_path)
        # ... actual training logic, e.g. calling into MLflow ...
        accuracy = 0.94
        return {"accuracy": accuracy, "model_path": "/tmp/model.pkl"}

    @task
    def deploy_if_good(metrics: dict):
        if metrics["accuracy"] < 0.90:
            raise ValueError(f"Accuracy {metrics['accuracy']} below deployment threshold")
        print(f"Deploying model from {metrics['model_path']}")

    raw_path = fetch_new_data()
    validated_path = validate_data(raw_path)
    metrics = train_model(validated_path)
    deploy_if_good(metrics)

ml_retraining_pipeline()

Notice you never explicitly wrote >> dependency arrows here — passing raw_path into validate_data(), and validated_path into train_model(), is exactly how the TaskFlow API infers task order and handles the underlying XCom data-passing automatically. This is genuinely the cleaner successor to manually chaining PythonOperator instances the older way.

Connecting This to MLOps: Where Each Piece Comes From

Recall the MLOps beginners guide's five-stage lifecycle — this DAG is quite literally that lifecycle expressed as executable code, and each task genuinely calls out to the specialized tools that article covered rather than reinventing them.

fetch_new_data — this is the ETL pipeline article's extract step, scheduled and monitored instead of run manually.

train_model — in a real deployment, this task would log to MLflow exactly as shown in the MLOps article's five-line tracking example, not just print a bare accuracy number.

deploy_if_good — this connects to the model serving layer (BentoML, KServe) from the MLOps platforms comparison, triggered conditionally based on the evaluation gate.

Airflow's actual job here is the glue, not the ML logic itself. It doesn't replace MLflow's experiment tracking or a dedicated serving framework — it's the scheduler and dependency manager coordinating when each of those tools' actual work happens.

Why "Retry Just the Failed Task" Matters So Much for ML Specifically

Here's a genuinely practical reason this matters more for ML pipelines than plain ETL: training runs can be expensive and slow. If your extract step fails at 3am due to a flaky API, Airflow sends an alert and lets you re-trigger just that failed task once the issue is fixed — not the entire pipeline, meaning you don't burn another multi-hour training run just to recover from an unrelated network hiccup upstream.

A real reported result: moving a manual reporting process to an Airflow DAG cut a retail client's reporting delay from 24 hours to 45 minutes, while eliminating the 3-4 manual steps someone had been doing by hand every morning. The equivalent ML story is the same idea applied to retraining — instead of someone remembering to kick off a retrain job weekly, the DAG handles it on schedule, fails loudly and specifically when something's wrong, and recovers without a full restart.

Critical Best Practices (Genuinely Non-Negotiable Ones)

A few of these aren't stylistic preferences — they affect whether your pipeline works correctly at all.

Don't put meaningful logic in top-level DAG code. Airflow's scheduler re-parses every DAG file at a fixed interval, independent of whether the DAG is even running — code sitting outside a task or operator executes on every single one of those parses, not just when the pipeline runs. This is a genuinely common, confusing source of unexpected behavior and real performance drag at scale.

Treat each task like a database transaction — never produce incomplete results. Don't write half a file to S3 and let a downstream task naively assume it's complete; design tasks to either fully succeed or leave no partial trace at all.

Don't pass large data through XComs. They're designed for small messages — a model's accuracy score, a file path, a run ID — not an entire dataset. For genuinely large data moving between tasks, write to shared storage (S3, a database) and pass the reference through XCom instead.

Never hardcode credentials inside a task. Use Airflow's Connections feature to store authentication details securely in the backend, retrieved by a connection ID rather than a literal password sitting in your DAG file.

Version control your DAGs like any other code. Defining workflows in Python specifically enables this — track changes in Git, roll back a bad DAG revision, and collaborate on pipeline changes the same way you would application code.

Event-Driven Scheduling: Beyond the Cron Schedule

Airflow 3.0 introduced Dataset-aware scheduling — instead of purely time-based triggers, a DAG can now run specifically when new data actually arrives, not just on a fixed clock.

from airflow.datasets import Dataset

@dag(schedule=[Dataset("s3://ml-bucket/raw-training-data/")])
def retrain_on_new_data():
    ...

This is genuinely more appropriate for ML retraining than a blind daily schedule in many cases — if new labeled data lands irregularly, triggering retraining specifically when it arrives is more sensible than either retraining on stale data daily or manually watching for new files yourself.

When Airflow Is Genuinely the Wrong Tool

Recall the ETL pipeline article's honest caveat, worth repeating specifically for ML: Airflow is not a library — you have to actually deploy and operate it, and that overhead isn't justified for a single-step job or infrequent, simple retraining that a cron job and a plain Python script handle just fine.

Airflow works best for workflows that are mostly static and slowly changing — the DAG structure itself should look similar run to run. If your pipeline's shape genuinely changes every single execution, Airflow's model fights you rather than helping.

It's explicitly not a streaming solution — it's built for batch-oriented, scheduled or event-triggered work. If you need genuine real-time stream processing, that's Kafka-plus-Spark territory from the data engineering courses article, not Airflow alone.

For high-volume, data-intensive tasks specifically, delegate the actual heavy lifting to specialized external services — Airflow orchestrates when a Spark job runs, it generally shouldn't be the thing crunching 25 million rows itself.

Common Mistakes People Make

Writing meaningful logic outside of task/operator boundaries. This runs on every scheduler parse cycle, not just during actual DAG execution — a genuinely confusing and avoidable performance and correctness trap.

Passing large datasets through XComs instead of shared storage references. XComs are for small metadata, not your actual training data — treat this boundary seriously.

Deploying Airflow for a single, infrequent job. The operational overhead of running Airflow itself outweighs the benefit for genuinely simple, small-scale pipelines — plain Python plus cron remains the right call there.

Skipping the data validation step before training. A silent bad-data problem upstream (recall the validate_data task above) can otherwise train and even deploy a genuinely broken model without any pipeline failure ever firing.

Hardcoding secrets directly in DAG files. Use Airflow Connections — credentials committed into a DAG file are a real, avoidable security liability.

  • Data Pipelines Pocket Reference by James Densmore — a concise, practical guide to designing and maintaining data pipelines, covering exactly the patterns Airflow automates. Perfect for understanding the "why" behind DAG design decisions.
  • Programming Amazon Web Services by James Murty & Ian Eーション — essential reading for the storage and compute services your Airflow tasks will actually orchestrate, from S3 data passes to EC2-based training jobs.
  • Designing Machine Learning Systems by Chip Huyen — the definitive guide to exactly what this article covers: getting ML models from notebook to production, including orchestration and pipeline design.
  • Machine Learning Engineering by Andriy Burkov — a rigorous, practical reference for the MLOps lifecycle from a practitioner's perspective, including pipeline orchestration patterns.

Wrapping This Up

Airflow orchestrates ML workflows the same way it orchestrates any other pipeline — DAGs define what depends on what, tasks do the actual work, and the scheduler handles retries, alerting, and recovery so a 3am failure doesn't require someone waking up to manually restart an entire multi-hour training run. The TaskFlow API's @task decorators are genuinely the cleanest current way to write these pipelines, inferring dependencies from how data actually flows between functions rather than requiring you to wire up explicit >> chains by hand.

Remember that top-level DAG code runs on every scheduler parse regardless of whether the pipeline is executing, and that XComs are for small metadata, not your actual datasets. FYI, this genuinely closes the loop on the MLOps arc from earlier in this series — Airflow is the orchestration layer coordinating when MLflow logs a run, when a model gets evaluated against a deployment threshold, and when BentoML or KServe actually serves the result, rather than replacing any of those tools itself :)

Now go take the plain-Python ETL pipeline from the previous article and wrap its three functions in @task decorators inside a DAG exactly like the example above. Watching the same logic gain scheduling, retries, and a monitoring UI for free is genuinely the clearest way to feel what Airflow actually adds.

Share this article X Facebook LinkedIn Reddit WhatsApp

Related Articles