A beginner’s introduction to Apache Airflow
What Apache Airflow is, where it comes from, what it is used for, and the core concepts of Airflow 3 to grasp before going further.
Introduction
This article is an introduction to Apache Airflow, for anyone who wants to understand its basic concepts. It covers what Airflow is, where it comes from, what it can be used for, and the concepts to grasp before going further.
Apache Airflow is an open-source tool to author, schedule and monitor workflows. It started at Airbnb in 2014, joined the Apache Software Foundation in 2016, and has been developed by the community ever since. The latest major version is 3.x, and this article describes the concepts of that version.
Airflow is a versatile workflow orchestrator with a wide range of use cases: from infrastructure management (spinning resources up and tearing them down as needed) to data pipelines (Extract Transform Load, ETL, or Extract Load Transform, ELT). It fits any process made of steps that run one after the other or in parallel.
A major benefit of Airflow is that workflows are defined as code. Software practices such as versioning and CI/CD pipelines apply to workflows too. Airflow also comes with a web interface to trigger workflows manually and to inspect their runs and logs when debugging.
Airflow concepts
Authoring and managing workflows is at the core of Apache Airflow. A workflow is a collection of tasks executed in a specific order. Airflow represents workflows as Directed Acyclic Graphs (DAGs): directed because tasks depend on each other and run in a given order, acyclic because no loop is allowed.
What is a DAG?
A DAG is the model of a workflow in Airflow, with all the metadata attached to it:
- the tasks of the workflow;
- the dependencies between tasks (the execution order);
- the schedule of the workflow runs;
- the start and end dates of the workflow.
A DAG describes which tasks make up the workflow, but doesn’t do the actual work: that is the job of the tasks. There are three ways to declare a DAG.
Method 1: context manager (recommended)
from datetime import datetime
from airflow.sdk import DAG
with DAG(
dag_id="my_dag",
start_date=datetime(2024, 1, 1),
schedule="@daily",
catchup=False,
) as dag:
# Tasks defined here are automatically assigned to this DAG
passMethod 2: DAG class constructor
from datetime import datetime
from airflow.sdk import DAG
dag = DAG(
dag_id="my_dag",
start_date=datetime(2024, 1, 1),
schedule="@daily",
catchup=False,
)
# Tasks must reference the DAG explicitly
# task = SomeOperator(..., dag=dag)Method 3: TaskFlow decorator
from datetime import datetime
from airflow.sdk import dag
@dag(
dag_id="my_dag",
start_date=datetime(2024, 1, 1),
schedule="@daily",
catchup=False,
)
def my_workflow():
# Define tasks with the @task decorator here
pass
my_workflow()Tasks and operators
Tasks are the unit of execution of a workflow: this is where the actual work happens. A task is an instance of an operator, the class that defines what the task does. Airflow ships with operators such as BashOperator and PythonOperator, which run Bash commands and Python functions respectively. In Airflow 3 they come from the standard provider, installed with Airflow.
BashOperator and PythonOperator only go so far for workflows that integrate with external services such as cloud providers or data platforms. This is where providers come in handy (see below).
Method 1: context manager
from datetime import datetime
from airflow.providers.standard.operators.bash import BashOperator
from airflow.providers.standard.operators.python import PythonOperator
from airflow.sdk import DAG
def my_python_function():
print("Hello from Python!")
with DAG(
dag_id="my_dag",
start_date=datetime(2024, 1, 1),
schedule="@daily",
) as dag:
# Tasks are automatically assigned to the DAG
task1 = BashOperator(
task_id="bash_task",
bash_command="echo 'Hello from Bash!'",
)
task2 = PythonOperator(
task_id="python_task",
python_callable=my_python_function,
)
# Define task dependencies
task1 >> task2Method 2: DAG class constructor
from datetime import datetime
from airflow.providers.standard.operators.bash import BashOperator
from airflow.providers.standard.operators.python import PythonOperator
from airflow.sdk import DAG
def my_python_function():
print("Hello from Python!")
dag = DAG(
dag_id="my_dag",
start_date=datetime(2024, 1, 1),
schedule="@daily",
)
# Tasks must reference the DAG explicitly
task1 = BashOperator(
task_id="bash_task",
bash_command="echo 'Hello from Bash!'",
dag=dag,
)
task2 = PythonOperator(
task_id="python_task",
python_callable=my_python_function,
dag=dag,
)
# Define task dependencies
task1 >> task2Method 3: TaskFlow decorator
from datetime import datetime
from airflow.sdk import dag, task
@dag(
dag_id="my_dag",
start_date=datetime(2024, 1, 1),
schedule="@daily",
)
def my_workflow():
@task
def extract():
return {"data": [1, 2, 3]}
@task
def transform(data: dict):
return {"transformed": [x * 2 for x in data["data"]]}
@task
def load(data: dict):
print(f"Loading: {data}")
# TaskFlow infers dependencies from the function calls
raw_data = extract()
transformed_data = transform(raw_data)
load(transformed_data)
my_workflow()Providers
Providers extend the core capabilities of Apache Airflow. They bring operators, sensors and other building blocks that make it easier to integrate Airflow with external systems.
Providers are installed separately from Airflow core, as Python packages. Airflow detects what a provider offers when it restarts after the provider is installed.
Airflow architecture
Apache Airflow is made of several components. It can be deployed as a distributed system (recommended for production) or on a single machine. These are the components that make it work.
Scheduler
The scheduler is the brain of Airflow. It reads the DAGs from the metadata database, triggers their runs when they are due, and hands their tasks to the executor, which runs inside the scheduler process.
API server
Serves the web interface and the REST API used to inspect, debug and trigger DAGs and tasks.
DAG processor
Continuously scans the DAG folder, then parses new and changed DAG files and serialises them into the metadata database.
Metadata database
The metadata database is the source of truth of the whole system: it stores the state of DAGs and tasks. Airflow supports PostgreSQL, MySQL and SQLite (not for production).
Conclusion
Apache Airflow is a versatile workflow orchestrator with many use cases. Workflows are defined as code, and a web interface helps inspect, debug and trigger DAGs and tasks.
For more advanced topics, head over to the official documentation.