Apache Airflow

Apache Airflow

Definition: Apache Airflow is an open-source platform for programmatically authoring, scheduling, and monitoring data pipelines, where a workflow is expressed as a DAG (Directed Acyclic Graph) of tasks written in Python. It was created at Airbnb in 2014 to replace a growing tangle of cron jobs coordinating the company’s internal data workflows, and was donated to the Apache Software Foundation two years later. Airflow does not move or transform data itself — it orchestrates other systems (Apache Spark clusters, databases, APIs, warehouses like Snowflake) that do the actual work, tracking what ran, what failed, and what depends on what. Its defining philosophy, “workflows as code,” means pipelines are ordinary Python, so anything Python can do — loops, conditionals, dynamic task generation — can shape the DAG itself.

Architecture

  • DAG (Directed Acyclic Graph): the core abstraction — a Python file declaring a set of tasks and the dependencies between them; “acyclic” guarantees no task can depend on itself, directly or transitively, so every run has a well-defined order
  • Operators: templates defining what a single task actually does — PythonOperator runs a Python function, BashOperator runs a shell command, PostgresOperator runs SQL, and provider packages add hundreds more for services like S3, BigQuery, or Snowflake
  • Tasks: an instantiated operator inside a specific DAG — the DAG is the blueprint, tasks are the concrete nodes actually scheduled and run
  • The Scheduler: a long-running process that continuously parses DAG files, determines which tasks are due based on schedule and dependencies, and queues them for execution
  • The Metadata Database: a relational database (commonly Postgres or MySQL in production) storing DAG state, task instance history, connections, variables, and XComs — effectively Airflow’s single source of truth
  • Executors: decide how and where queued tasks actually run — LocalExecutor runs tasks as subprocesses on the scheduler’s own machine, CeleryExecutor distributes them across a pool of workers via a message broker, KubernetesExecutor launches a fresh pod per task
  • Workers: the processes that pull tasks off a queue and execute the operator’s actual code — long-running in CeleryExecutor, ephemeral (one pod per task) in KubernetesExecutor
  • “DAGs as Python code”: the design philosophy that pipelines belong in a general-purpose language rather than YAML or a drag-and-drop UI, since loops and functions let a DAG be generated dynamically instead of hand-written task by task

Why It Matters

  • Became the de facto standard for pipeline orchestration, replacing sprawling collections of cron jobs with no shared visibility into failures or dependencies
  • The web UI gives one place to see run history, task duration trends, logs, and failures across every pipeline, instead of SSH-ing into a box to check a cron log
  • Retries, alerting, and SLAs are handled declaratively per task, so failure handling doesn’t have to be hand-rolled into every script
  • Being “just Python” means it integrates with virtually any system that has a client library or REST API, and skills engineers already have — functions, loops, version control — transfer directly
  • A large open-source ecosystem of provider packages means most common integrations already exist as pre-built operators rather than needing custom code

Under the Hood: The Scheduler Loop

  1. The scheduler continuously scans the configured DAGs folder, re-parsing each Python file on a set interval (commonly every 30 seconds) to pick up new or edited DAGs
  2. For each parsed DAG, it checks the schedule (a cron expression or preset like @daily) against the last recorded run in the Metadata Database to decide whether a new run is due
  3. When a run is due, the scheduler creates a DAG run record and evaluates every task’s dependencies — a task becomes eligible the moment all of its upstream tasks reach a satisfying state
  4. Eligible tasks are handed to the configured executor, which places them on a queue rather than running them itself — the scheduler’s job stops at “this task is ready”
  5. The executor dispatches queued tasks to available workers, which pick them up, execute the operator’s code, and stream logs back
  6. As each task finishes, the worker reports its final state (success, failed, up_for_retry) to the Metadata Database, which the scheduler reads on its next loop to decide what’s now eligible downstream
  7. This loop — parse, evaluate, queue, dispatch, record — repeats continuously, which is why a DAG’s actual trigger time is always “on or after” its schedule, never exactly on it

Comparison: Airflow vs Dagster vs Prefect

DimensionAirflowDagsterPrefect
Scheduling modelCentralized scheduler polling loop against a metadata DBCentralized daemon with tighter run coordinationHybrid — scheduler-driven or purely event/API-triggered
DAG authoring styleTasks and dependencies via operators and >> syntax“Software-defined assets” — declares data assets and lineage, not just tasksPlain Python functions decorated as tasks/flows, dependencies inferred from calls
Typical use caseBroad, general-purpose orchestration across heterogeneous systemsAsset-centric pipelines where lineage and testing are first-classLightweight, dynamic, Python-native and event-driven jobs
Operational complexityModerate to high — scheduler, metadata DB, executor, and workers to run and tuneModerate — similar footprint, own UI and daemonLower for simple deployments — can run without a persistent scheduler
Ecosystem maturityLargest, oldest, most provider packagesSmaller but fast-growing, strong typing and testing storySmaller, more recent, popular with lighter-weight teams

Common Pitfalls

  • Putting expensive top-level code (API calls, DB queries, large computations) directly in a DAG file — the scheduler re-parses every DAG file on every cycle, so slow top-level code slows down the whole scheduler, not just that DAG
  • Writing tasks that aren’t idempotent — a task that appends a row or increments a counter produces a different result on retry than on first run, breaking Airflow’s entire retry-and-backfill model
  • Passing large amounts of data between tasks via XComs, which are meant for small metadata values, not dataframes or files — pushing megabytes through XComs bloats the Metadata Database and slows every scheduler query against it
  • Undersizing the scheduler or executor for actual task volume — a scheduler starved of CPU falls behind on parsing and evaluation, which shows up as tasks stuck in a queued state for no visible reason
  • Treating Airflow as a data processing engine instead of an orchestrator — running heavy transformations inside a PythonOperator instead of delegating to Apache Spark or a warehouse’s own compute overloads workers never sized for that
  • Letting DAGs grow into deeply nested, hard-to-read dependency chains that defeat the original goal of replacing an unreadable tangle of cron jobs with something legible
  • Hardcoding credentials or environment-specific values in DAG files instead of using Connections and Variables, making DAGs impossible to promote cleanly between environments — see Infrastructure as Code (IaC)
  • Leaving catchup on by default for a newly deployed DAG with a start date months in the past, triggering a flood of backfill runs all at once

Worked Example

StepComponentWhat happens
1SchedulerDetects the DAG’s @daily interval has elapsed and creates a new DAG run
2SchedulerParses the DAG file, builds the dependency graph, marks extract_from_api and extract_from_db ready (no upstream dependencies)
3ExecutorQueues both extraction tasks; with no dependency between them, they dispatch concurrently
4WorkerRuns extract_from_api, writes its output location to XCom, reports success
5WorkerRuns extract_from_db in parallel, reports success
6SchedulerSees both upstream tasks succeeded, marks validate_schema then merge_datasets ready as dependencies clear
7Workertransform fails on first attempt due to a transient network error
8SchedulerSees the failure, checks configured retries, marks the task up_for_retry, and re-queues it after the retry delay
9Workertransform succeeds on retry, unblocking load_to_warehouse
10SchedulerMarks the DAG run successful once every task reaches a terminal success state, visible in the web UI

Real-World Use

  • ETL/ELT orchestration: coordinating extract, transform, and load steps across databases, APIs, and warehouses like Snowflake on a fixed schedule — see ETL vs ELT
  • ML pipeline scheduling: triggering feature engineering, model training, and batch scoring jobs in sequence, often gated by a model registry step before deployment
  • Cross-system workflow coordination: kicking off an Apache Spark job, waiting for it to finish, then triggering a downstream notification or another team’s pipeline via sensors
  • Data quality check scheduling: running validation suites after a load completes and blocking downstream consumers until checks pass — see Data Quality and Validation
  • Batch reporting pipelines: regenerating dashboards, exports, or scheduled reports on a nightly or hourly cadence with full visibility into which run produced which output

Best Practices

  • Keep DAG files lightweight — do heavy lifting inside task functions or external systems, not at import time, so parsing stays fast
  • Design every task to be idempotent — re-running it with the same inputs should produce the same end state, which is what makes retries and backfills safe
  • Use task groups to organize large DAGs visually and logically without changing the underlying dependency structure
  • Set sensible retries, retry_delay, and SLAs per task rather than relying on defaults, so transient failures self-heal and real failures alert promptly
  • Avoid passing large payloads through XComs — write intermediate data to object storage or a warehouse and pass a reference instead
  • Version-control DAGs like any other code, with review and CI checks that at minimum verify every DAG file parses without error

FAQ

Is Airflow a data processing engine? No — it’s an orchestrator. It decides what runs, when, and in what order, but the actual data movement or transformation should happen in the systems it calls out to, like Apache Spark or a warehouse.

Does Airflow move data between tasks? Not directly, beyond small values passed via XCom. Tasks typically read and write through external storage, with Airflow only coordinating when each step runs.

What happens if a scheduled run is missed because Airflow was down? Depending on the catchup setting, Airflow either backfills every missed interval on restart or skips straight to the next scheduled run — a common source of surprise flooding after a redeploy.

Can Airflow trigger DAGs outside of a fixed schedule? Yes — DAGs can be triggered manually, via the REST API, or by sensors and datasets that fire when an upstream condition, like another DAG finishing or a file landing, is met.

Is Airflow suitable for real-time or streaming pipelines? Not natively — it’s built around discrete, scheduled batch runs, not continuous processing; streaming workloads are usually better served by something like Apache Kafka, with Airflow orchestrating the surrounding batch steps — see Batch vs Stream Processing.

How does Airflow differ from a plain cron job? Cron only knows “run this at this time,” with no concept of dependencies, retries, history, or failure visibility; Airflow adds all of that plus a UI, at the cost of real operational overhead a single cron entry doesn’t have.

History

  • Created at Airbnb in 2014 by Maxime Beauchemin, to replace an unwieldy collection of cron jobs coordinating the company’s internal data pipelines
  • Open-sourced by Airbnb in mid-2015
  • Entered the Apache Software Foundation Incubator in 2016, and graduated to a Top-Level Apache project in 2019
  • Maxime Beauchemin later created Apache Superset and then Preset and Dagster, drawing directly on lessons learned building Airflow
  • Airflow 2.0, released in December 2020, was a major rewrite adding a high-availability scheduler, a stable REST API, and a redesigned UI
  • The “DAGs as Python code” model Airflow popularized directly shaped a wave of later orchestrators — Prefect, Dagster, and Luigi each responded to specific pain points teams hit running Airflow at scale
  • Later major versions continued modernizing the platform with DAG versioning and a stronger separation between the scheduler and task execution environments

Common Interview Questions

  • “What is a DAG, and why does it need to be acyclic?” — expect an explanation that cycles would make “what runs first” undefined, since a task could end up depending on its own output
  • “How would you handle a task that needs to pass a large dataset to the next task?” — expect a rejection of XComs for this purpose, in favor of writing to shared storage and passing a reference
  • “What’s the difference between the scheduler and the executor?” — expect a clear separation: the scheduler decides what’s ready, the executor decides where and how it actually runs
  • “How do you make a DAG idempotent?” — expect discussion of partition-based writes (overwrite a specific date’s partition rather than append) and avoiding side effects that depend on run count
  • “What would you check if tasks are stuck in a queued state?” — expect a mention of executor or worker capacity, scheduler resource starvation, or a misconfigured pool limiting concurrency

Example

An Airflow DAG scheduled for 6 AM daily starts by extracting yesterday’s orders from a REST API and, in parallel, pulling customer records from a production database. Once both extractions succeed, a validation task checks the API response against an expected schema before a merge task combines the two datasets. A transform task then cleans and reshapes the merged data, and a final task loads the result into a Snowflake warehouse table, overwriting only that day’s partition so the run can be safely retried or backfilled without duplicating data. If the transform step fails — say, due to a transient timeout — Airflow automatically retries it after a configured delay before alerting the on-call engineer, and the web UI shows exactly which task in the chain failed, how long each step took, and the full log output, without anyone needing to SSH into a box or grep through a cron log.

Dig deeper