Workflow Orchestration (Airflow & Prefect)
Learn how apache-airflow and prefect schedule and orchestrate multi-step data pipelines as code, and how to choose between them.
Introduction
A real data pipeline is rarely one script. It is usually a sequence of steps — pull raw data, clean it, train or update a model, publish a result — that need to run in order, on a schedule, with retries when something fails and alerts when it fails anyway. Running that by hand with a cron job works until it does not. This lesson covers the two most common Python libraries built specifically to orchestrate pipelines like this.
- How apache-airflow represents a pipeline as a DAG (directed acyclic graph) of tasks.
- How prefect offers a lighter, more Python-native way to define the same kind of workflow.
- How to choose between the two for a given project.
apache-airflow: DAG-Based Orchestration
Apache Airflow is the most widely adopted workflow orchestration tool in data engineering. A pipeline in Airflow is defined as a DAG — a directed acyclic graph of tasks, where each task depends on the ones before it, and the whole graph has no cycles. Airflow provides a scheduler that runs DAGs on a timer, a web UI to monitor runs, and built-in retry and alerting logic.
pip install apache-airflowfrom airflow import DAGfrom airflow.operators.python import PythonOperatorfrom datetime import datetime
def extract_data(): print("Pulling latest customer data from the warehouse...")
def train_model(): print("Retraining the churn model on fresh data...")
with DAG( dag_id="churn_pipeline", start_date=datetime(2026, 1, 1), schedule="@daily", catchup=False,) as dag:
extract_task = PythonOperator( task_id="extract_data", python_callable=extract_data, )
train_task = PythonOperator( task_id="train_model", python_callable=train_model, )
extract_task >> train_task # train_task runs only after extract_task succeedsThe line extract_task >> train_task declares a dependency: train_task will not start until extract_task finishes successfully. This is how Airflow builds the graph shape of a DAG out of individual task objects.
Click Run to see what this code prints.
prefect: Python-Native Orchestration
Prefect solves the same core problem as Airflow — scheduling and orchestrating multi-step pipelines — but with a lighter-weight, more Python-native feel. Instead of building a DAG object explicitly, you write regular Python functions and decorate them with @task and @flow. Prefect infers the dependency graph from how you actually call the functions.
pip install prefectfrom prefect import flow, task
@taskdef extract_data(): print("Pulling latest customer data from the warehouse...") return {"rows": 5000}
@taskdef train_model(data): print(f"Retraining the churn model on {data['rows']} rows...")
@flow(name="churn-pipeline")def churn_pipeline(): data = extract_data() train_model(data)
if __name__ == "__main__": churn_pipeline()Click Run to see what this code prints.
Notice that train_model(data) depends on extract_data() simply because it receives its output as an argument — there is no separate ">>" syntax to declare the graph. Prefect builds the dependency graph from ordinary Python function calls.
Airflow vs Prefect
| Apache Airflow | Prefect | |
|---|---|---|
| Pipeline definition | Explicit DAG objects and operators | Plain Python functions with @task/@flow |
| Learning curve | Steeper — its own concepts and setup | Gentler — feels like regular Python |
| Ecosystem maturity | Very mature, huge library of integrations | Newer, growing quickly |
| Best fit | Large, established data engineering teams | Smaller teams wanting to move fast |
Common Mistakes
- Reaching for Airflow or Prefect on a single script that runs once a day with no real dependency chain — a simple cron job is often enough.
- Putting heavy data processing logic directly inside a DAG file, instead of calling out to well-tested functions or scripts.
- Forgetting to set retries and alerting on tasks that call external, unreliable services like APIs or databases.
- Confusing schedule_interval/schedule (when a DAG runs) with catchup (whether missed past runs get backfilled) in Airflow.
Best Practices
- Keep individual tasks small and focused — one clear responsibility each.
- Set retries and failure alerts on any task that depends on a network call or external system.
- Start with Prefect if the team is small and wants a fast, Python-native ramp-up; consider Airflow when the org already runs it or needs its mature ecosystem of integrations.
- Version-control pipeline definitions the same as any other code, with tests for the underlying task functions.
Frequently Asked Questions
Usually not. A simple scheduled script (via cron or a basic while loop) is enough until you have multiple dependent steps, need retries and monitoring, or are coordinating pipelines across a team.
For most day-to-day orchestration needs, yes. Airflow still has a larger, more mature ecosystem of pre-built integrations (called providers), which matters more at large-scale data engineering organizations.
It is most associated with data engineering, but a DAG is a general concept — Airflow can orchestrate any sequence of scheduled, dependent tasks, including ML training pipelines like the one in this lesson.
Key Takeaways
- apache-airflow orchestrates pipelines as explicit DAGs of tasks with a mature scheduler and UI.
- prefect orchestrates the same kind of pipeline using plain Python functions decorated with @task and @flow.
- Both add retries, scheduling, and monitoring that a plain script does not have on its own.
- Reach for either once a pipeline has multiple dependent steps that need to run reliably and on a schedule.
Summary
Airflow and Prefect turn a chain of scripts into a reliable, scheduled, monitored pipeline. Next, we look at how to track what actually happened during those pipeline runs — the experiments, parameters, and metrics — with MLflow and Weights & Biases.
- You can define a two-task DAG in Airflow.
- You can define an equivalent flow in Prefect using @task and @flow.
- You can choose between them for a given team and project.