Tutorial4 min read

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.

AirflowPythonOrchestration

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
    pass

Method 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 >> task2

Method 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 >> task2

Method 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.

Becko Junior Camara

DevOps Engineer (Azure & AWS)