1 of 28

Introduction to Apache Airflow

Xia Wang (王夏)

2021-12-14

2 of 28

  • What Is Airflow
  • When (not) to Use Airflow
  • Airflow’s Concepts/ Example
  • Airflow’s Architecture
  • Q&A

3 of 28

  • What Is Airflow
  • When (not) to Use Airflow
  • Airflow’s Concepts/ Example
  • Airflow’s Architecture
  • Q&A

4 of 28

Apache Airflow

  • “Airflow is a platform to programmatically author, schedule and monitor workflows.”
  • Created in 2014 in Airbnb
  • Open sourced in 2015
  • Joined Apache Software Foundation in 2016 and became a top level project in 2019

5 of 28

Apache Airflow

  • Workflow as Code
  • Use Python code to
    • Orchestrate tasks (establish task dependencies)
    • Schedule workflows

6 of 28

  • What Is Airflow
  • When (not) to Use Airflow
  • Airflow’s Concepts/ Example
  • Airflow’s Architecture
  • Q&A

7 of 28

Good at:

  • Batched Extract, Transform, Load (ETL) workflows (non-circular)

  • Scheduling regularly repeated workflows
    • Use cron expression or Python timedelta by default
    • Minute-level frequency or less frequent jobs
  • Tasks with dependent relationships
  • Static workflow structure

Extract frontend events (E)

Transform to session level (T)

Load session info to data warehouse (L)

8 of 28

Not Good at:

  • Data Streaming
  • Tasks with circular dependencies
  • An overwhelming number of workflows/ tasks
    • 1 workflow with 10k+ tasks
    • thousands of workflows each with hundreds of tasks

A

B

C

9 of 28

Airflow

Luigi

Oozie

Azkaban

Jenkins

workflow definition

Python

Python

XML

Custom DSL

Groovy/Bash

community

Very Active

Active

Active

Somewhat Active

Active

main purpose

general purpose Batch ETL

general purpose Batch ETL

Hadoop Job Scheduling

Hadoop Job Scheduling

CI/CD

scheduling

easy to schedule with the scheduler component

no central process that automatically triggers jobs

verbose scheduling definition

easy to schedule

cronjobs

task dependency/ parallelisation

out-of-box support for task dependency/ parallelisation

no straightforward way to make sure second task starts before first task ends

yes in Hadoop ecosystem

yes in Hadoop ecosystem

use stages to set execution order; need to manually define more complex workflows

easiness to use

easy

easy

relatively easy

relatively easy

hard

target audience

data engineer/analyst/scientist

data engineer/analysts/scientist

data engineer

data engineer

system operation/ dev team

10 of 28

  • What Is Airflow
  • When (not) to Use Airflow
  • Airflow’s Concepts/ Example
  • Airflow’s Architecture
  • Q&A

11 of 28

Basic Concepts in Airflow

  • Task: basic building block.
    • Task_run/task_instance (one instantiated task)
    • Operator: task templates
    • Sensor: task that needs to wait for an external event
    • TaskFlow decoratored @task: customized python functions

12 of 28

By default, Airflow supports, among others:

  • Python Operator (to execute Python functions)
  • Bash Operator (to execute Bash scripts or trigger a command/submit job, etc)
  • Postgresql/MySQL Operator (to execute a sql script in a DB)

Strong support from the community for a good variety of cloud service providers:

  • AWS (Redshift Operator, DMS Operator, EC2 instance Sensor, etc)
  • GCP (BigQuery, etc)
  • Azure (Cosmos, etc)
  • Slack (for sending slack messages)
  • A lot more

13 of 28

  • Directed Acyclic Graph (dag): built by tasks and dependencies
    • Dag_run (one instantiated dag)

14 of 28

  • Schedule Interval: how often a dag is executed (minute/hour/day/month/day of the week)
    • Cron expression
    • 0 5 * * *
    • 30 */2 * * *

15 of 28

Link to the dag class parameter reference

16 of 28

  • What Is Airflow
  • When (not) to Use Airflow
  • Airflow’s Concepts/ Example
  • Airflow’s Architecture
  • Q&A

17 of 28

18 of 28

  • Directory where the dag files are held
    • ~/airflow/dag/ for example
  • Metadata DB
    • recommended PG 9.6+ or MySQL 8+
  • Webserver
    • Provide UI
    • Visualize dag dependency
    • Stateless
  • Scheduler
    • Parses the dags into dag serialization and store state in DB
    • Schedule dags for execution (through executor)
  • Executor
    • Kubernetes
    • Celery
    • A combination of both
  • Worker: where task execution actually happens

19 of 28

More on Scheduler and HA

  • Parses the dags frequently (1 minute by default) to create dag serialization/trigger dag_run
  • Can have multiple schedulers running since 2.0
    • Better performance
    • More resilience
  • All schedulers communicate with the DB directly
    • Row level lock
    • Table level lock
      • 50-100 ms per scheduler per time

20 of 28

21 of 28

22 of 28

Kubernetes Executor: link

23 of 28

24 of 28

Celery Executor: link

25 of 28

26 of 28

27 of 28

Q & A

28 of 28

Reference

Airflow’s official documentation: link

Airflow vs Luigi comparison: link1 link2

Airflow vs Azkaban vs Oozie vs some other alternatives: link

Airflow’s helm chart installation: link

Deep dive in the Airflow Scheduler video