skip to content

Apache Airflow

The de-facto orchestrator: DAGs written in Python, operators and sensors for the work, a scheduler driving intervals, and pluggable executors. Interviewers focus on Airflow because its scheduling semantics and XCom limits shape how pipelines get designed.

on this pageshow

explore

questions

page 1 of 2

In Airflow, what do the >> and << operators do between two tasks in a DAG?

level: juniorimportance: must knowfreq 82%

answer

  1. the arrow is an edge, not a pipe
  2. it says after, not immediately after
  3. same thing as set_downstream
  4. lists broadcast, list-to-list does not
  5. order is decided at parse time

basics

~10 s

In Airflow, a >> b puts b downstream of a: b is not started until a has finished successfully. The arrow declares execution order only. It moves no data between the tasks.

solid answer

~40 s

`a >> b` is shorthand for `a.set_downstream(b)`, and `b << a` writes the same edge the other way. It adds a directed edge that the scheduler uses to decide eligibility: with the default `trigger_rule='all_success'`, `b` becomes schedulable only once every upstream task has succeeded. Lists give you fan-out and fan-in in one line — `extract >> [clean_a, clean_b] >> load` — but list-to-list is not supported by the operators; use the `chain()` helper for that. Two things the arrow does **not** do: it does not pass data (that is XComs or shared storage), and it does not promise `b` runs immediately or on the same worker — only that it runs after. Edges are built when the DAG file is parsed, so a task cannot add one at run time.

code

python · 7 lines
python
extract = PythonOperator(task_id="extract", python_callable=pull_orders)
clean_a = PythonOperator(task_id="clean_orders", python_callable=clean_orders)
clean_b = PythonOperator(task_id="clean_customers", python_callable=clean_customers)
load = PythonOperator(task_id="load", python_callable=load_warehouse)

# fan-out then fan-in: load waits for BOTH cleaners
extract >> [clean_a, clean_b] >> load

go deeper

for a junior

Be able to read and write both forms fluently and to sketch the resulting graph on a whiteboard, including the fan-out and fan-in forms with a list on one side.

for a middle

Explain that the edge only encodes ordering evaluated against a trigger rule, that it carries no data, and that the graph is fixed when the file is parsed rather than while a run executes.

for a senior

Point out the operational consequences: tasks joined by an arrow can run minutes apart on different workers, so shared local state is a bug, and run-time-dependent shape needs branching or dynamic mapping instead.

for a principal

Push for a house style — one direction, helpers like chain for long sequences, dependencies declared in one readable block — because a graph that reviewers cannot read at a glance is where incorrect ordering hides.

## What the arrow declares An Airflow DAG is a directed acyclic graph whose nodes are tasks and whose edges are ordering constraints. The bitshift operators are Airflow's syntax for adding those edges. `a >> b` reads "a then b" and is exactly equivalent to `a.set_downstream(b)`; `b << a` is the mirrored form, equivalent to `b.set_upstream(a)`. Airflow implements this by overriding `__rshift__` and `__lshift__` on `BaseOperator` (and on the `XComArg` objects the TaskFlow `@task` decorator returns), so the syntax works on operator instances, on lists of them, and on `TaskGroup` objects. The edge is graph metadata, nothing more. When a DAG run is created, every task instance starts in the `none` state and the scheduler repeatedly asks whether each task's dependencies are met. The default answer rule is `trigger_rule='all_success'`: all direct upstream tasks must be in `success` before the task is queued. Change the trigger rule and the same edge is still there but the condition on it changes — the edge says *ordering*, the trigger rule says *under what upstream outcome*. ## The four ways to write the same thing ```python extract >> transform # bitshift, the idiomatic form transform << extract # identical edge, reversed reading extract.set_downstream(transform) transform.set_upstream(extract) ``` All four produce the same serialized DAG. Teams pick one style and stay with it; mixing `>>` and `<<` in one file is a readability complaint reviewers raise often, because the reader has to re-establish direction on every line. ## Fan-out, fan-in, and the list rule A list on either side broadcasts: ```python extract >> [clean_orders, clean_customers] >> load ``` That is four edges: extract to each cleaner, and each cleaner to load. `load` waits for **both** cleaners under the default trigger rule. What does *not* work is list on both sides — `[a, b] >> [c, d]` raises a Python error, because a plain list has no `>>` behaviour to fall back on. When you genuinely want that, use the helpers: `chain(a, b, c)` wires a sequence, `chain([a, b], [c, d])` pairs same-length lists element-wise, and `cross_downstream([a, b], [c, d])` builds the full cross product. ## What the arrow is not **It is not data flow.** `extract >> transform` gives `transform` no access to whatever `extract` computed. Values travel through XComs (small metadata) or, for anything of size, through object storage with only the path passed along. Candidates who say "the arrow passes the return value" are conflating the dependency with TaskFlow's implicit XCom, which is a *separate* mechanism that happens to create the edge as a side effect. **It is not "runs immediately after".** Once dependencies are met the task is queued; when it actually starts depends on the executor, free slots, pool capacity and DAG-level concurrency limits. Two tasks joined by an arrow routinely run on different workers, in different processes, possibly minutes apart. Anything that relies on shared local state — a file written to `/tmp`, an open database session, a Python variable — breaks the moment the two tasks land on different machines. **It is not dynamic.** Dependencies are declared while the DAG file is being parsed, before any run exists. Writing `if result > 10: a >> b` inside a task's Python body has no effect on the graph; the arrow ran on a worker, long after the scheduler built its picture of the DAG. Run-time shape changes need the mechanisms designed for it: branching operators to *skip* paths, or dynamic task mapping to fan out over a list that is only known at run time. ## Cycles The A in DAG is enforced. If your declarations create a loop, Airflow raises a cycle error when the file is parsed and the DAG is marked as broken rather than being scheduled with an arbitrary order. The most common accidental cycle comes from a loop that wires each element to the next without breaking on the last item, or from wiring a TaskGroup back into one of its own members. ## Reading a graph in an interview Be fluent at translating both directions. Given `a >> [b, c] >> d >> e`, you should be able to say instantly: a runs first; b and c run in parallel once a succeeds; d waits for both; e waits for d. And given that description you should be able to write the two lines back.

  • If `a >> b` does not move data, how does `b` get a value that `a` computed?
    Through XCom: `a` pushes a small value into the metadata database and `b` pulls it, either explicitly with `xcom_pull` or implicitly by taking a TaskFlow function's return value as an argument. XCom is sized for identifiers and small metadata, not datasets — for anything large, write to object storage in `a` and pass only the path.
  • What happens if the declarations in a DAG file create a cycle?
    Airflow detects it while parsing the file and raises a cycle error, so the DAG shows as broken with an import error instead of being scheduled. There is no partial run and no arbitrary tie-break: acyclicity is a precondition for the scheduler being able to answer "are this task's dependencies met" at all.
  • Why does `[a, b] >> [c, d]` fail, and what do you use instead?
    The left operand is a plain Python list, which has no bitshift behaviour, so there is nothing for Airflow to hook into. Use `chain([a, b], [c, d])` to pair the two lists element-wise, or `cross_downstream([a, b], [c, d])` to connect every task on the left to every task on the right.

saying these in an interview costs you the question

  • Says the arrow passes the upstream task's return value
  • Claims the downstream task starts immediately after the upstream one
  • Assumes both tasks share a worker, filesystem or process
  • Adds dependencies inside a task body and expects the graph to change
  • Thinks >> and << mean different kinds of dependency

context

open as a page

In Airflow, what do the task-level retries and retry_delay arguments control?

level: juniorimportance: must knowfreq 78%

basics

~10 s

In Airflow, retries sets how many extra attempts a task instance gets after its first failure, and retry_delay is the wait between attempts. Between attempts the task sits in the up_for_retry state.

open as a page

In Apache Airflow, what is the difference between an operator and a sensor?

level: juniorimportance: must knowfreq 80%

basics

~20 s

In Airflow an operator is a template for one unit of work — run a Bash command, call a Python function, submit a query. A sensor is an operator that only waits, polling until an external condition becomes true.

open as a page

In Airflow, what does setting catchup=False on a DAG do?

level: juniorimportance: must knowfreq 78%

basics

~20 s

With catchup=False, Airflow schedules only the most recent completed data interval when a DAG becomes active, instead of creating one run for every interval back to start_date. Older intervals are skipped and must be backfilled deliberately.

open as a page

In Airflow, what is an XCom and how does one task push a value another task pulls?

level: juniorimportance: must knowfreq 80%

basics

~10 s

An XCom is Airflow's small cross-task message. A task pushes a keyed value into the metadata database, and a downstream task in the same DAG run pulls it back by task id and key.

open as a page

In Airflow's TaskFlow API, how does calling one @task function inside another build the DAG?

level: middleimportance: must knowfreq 70%

basics

~20 s

Calling a @task-decorated function does not run it. It returns a placeholder for that task's future output, and passing that placeholder into another @task call both creates the dependency edge and arranges an XCom pull at run time.

open as a page

In Airflow, which settings limit how many task instances can run at the same time?

level: middleimportance: must knowfreq 62%

basics

~10 s

Airflow throttles at several levels: [core] parallelism per scheduler, max_active_runs and max_active_tasks per DAG, max_active_tis_per_dag per task, pool slots across DAGs, and the worker capacity your executor provides. The tightest one wins.

open as a page

In Airflow, how do LocalExecutor, CeleryExecutor and KubernetesExecutor differ in where tasks run?

level: middleimportance: must knowfreq 80%

basics

~10 s

LocalExecutor runs tasks as subprocesses on the scheduler's machine, CeleryExecutor sends them through a broker to separate long-lived worker processes on other machines, and KubernetesExecutor creates one short-lived pod per task instance.

open as a page

How do you get alerted when an Airflow task fails, and what do the callback arguments give you?

level: middleimportance: must knowfreq 72%

basics

~20 s

Airflow fires per-task callables such as on_failure_callback, on_retry_callback and on_success_callback, plus email_on_failure to a configured SMTP address. Callbacks receive a context dictionary with the DAG, task, run and log URL, which you forward to Slack, PagerDuty or a webhook.

open as a page

In an Airflow sensor, what is the difference between poke mode and reschedule mode?

level: middleimportance: must knowfreq 70%

basics

~20 s

Poke mode keeps the sensor task running for the whole wait, sleeping between checks and holding a worker slot. Reschedule mode ends the task after each failed check and re-queues it, freeing the slot until the next check is due.

open as a page

In Airflow, why does a DAG run's data_interval_start sit behind the time the run actually starts?

level: middleimportance: must knowfreq 82%

basics

~20 s

Airflow names a run for the data window it covers, not the moment it fires. A daily DAG's run whose data_interval_start is May 1 only starts once that day has closed, at May 2 midnight, so the data is complete.

open as a page

Which Airflow task instance states do you check first when a DAG run appears stuck?

level: seniorimportance: must knowfreq 66%

basics

~20 s

Read the task instance state in Airflow's grid view: scheduled or queued means no worker slot, running means the work is live, up_for_retry means it is waiting between attempts, and none or upstream_failed means it was never eligible. The state names the layer to investigate.

open as a page

An Airflow task returns a 2 GB DataFrame via XCom — what breaks, and what should it pass instead?

level: seniorimportance: must knowfreq 66%

basics

~20 s

It fails: a DataFrame is not JSON-serializable, and even serialized it would not fit the metadata database's XCom column. Write the frame to object storage or a table and push only the URI or partition key.

open as a page

In Airflow, how does declaring a DAG with the @dag decorator differ from the DAG() context manager?

level: juniorimportance: should knowfreq 60%

basics

~20 s

Both produce the same DAG. The DAG() context manager registers every operator built inside the with block; the @dag decorator turns a function into a DAG factory that you must call at module level, or nothing is registered.

open as a page

In Airflow, which long-running processes make up a deployment and what does each do?

level: juniorimportance: should knowfreq 62%

basics

~20 s

An Airflow deployment runs a scheduler that decides what should run, workers that execute task code, a web UI/API process, a triggerer for deferred tasks, and a metadata database that all of them read and write.

open as a page

In Airflow, what does calling .expand() on a task produce when the DAG runs?

level: middleimportance: should knowfreq 52%

basics

~20 s

It produces one mapped task instance per element of the expanded input, created at run time once the upstream value is known. The DAG shows a single task node; the run shows N instances, each identified by its map_index.

open as a page

In Airflow, how does KubernetesExecutor differ from using KubernetesPodOperator?

level: middleimportance: should knowfreq 48%

basics

~20 s

KubernetesExecutor is a deployment-wide choice that runs every Airflow task in its own worker pod. KubernetesPodOperator is a single operator, usable under any executor, that launches an arbitrary pod of your choosing and waits for it.

open as a page

In Airflow, what does setting an sla on a task do, and how does it differ from execution_timeout?

level: middleimportance: should knowfreq 54%

basics

~20 s

In Airflow an sla is a timedelta measured from the DAG run's data interval end; if the task has not succeeded by then, Airflow records an SLA miss and fires sla_miss_callback. It never interrupts the task, unlike execution_timeout, which kills a running attempt.

open as a page

In Airflow, what happens to the downstream tasks a BranchPythonOperator does not return?

level: middleimportance: should knowfreq 52%

basics

~20 s

They are marked skipped. The callable returns the task_id or ids to follow; every other directly downstream task is skipped, and because the default trigger rule requires successful upstreams, that skip cascades down the whole unchosen branch.

open as a page

In Airflow, why use a provider operator or hook instead of calling boto3 in a PythonOperator?

level: middleimportance: should knowfreq 44%

basics

~20 s

Provider operators and hooks resolve credentials through Airflow Connections and the secrets backend, template their arguments, log and retry consistently, and expose parameters in the UI. Hand-rolled client code usually re-implements that badly and hardcodes secrets.

open as a page

What replaced execution_date and schedule_interval in Airflow 3?

level: middleimportance: should knowfreq 50%

basics

~20 s

Airflow 3 removed execution_date in favour of logical_date plus the data_interval_start and data_interval_end pair, and replaced schedule_interval and the separate timetable argument with one schedule argument that accepts a cron string, a timedelta, a Timetable, or assets.

open as a page

Where does Airflow store XCom values by default, and what limits their size?

level: middleimportance: should knowfreq 70%

basics

~20 s

Airflow serializes each XCom to JSON and writes it as a row in the xcom table of its metadata database. The ceiling is that column's type, so the practical budget is kilobytes to a few megabytes, not gigabytes.

open as a page

In Airflow's TaskFlow API, how does returning a value from an @task function create an XCom?

level: middleimportance: should knowfreq 62%

basics

~20 s

The @task decorator wraps the function in an operator whose return value is pushed as the return_value XCom. Calling the function in a DAG returns a lazy reference, and passing it to another task both wires the dependency and generates the pull.

open as a page

In Airflow, why does top-level code in a DAG file slow down the whole scheduler?

level: seniorimportance: should knowfreq 58%

basics

~20 s

Because everything at module level runs every time the file is parsed, and Airflow re-parses every DAG file on a short repeating interval. An API call or query at the top level is paid on every parse, for every file, delaying scheduling everywhere.

open as a page

Airflow task instances sit in the queued state for hours and never start — how do you diagnose it?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Queued means Airflow's scheduler handed the task to the executor but nothing picked it up. Check in order: pool and concurrency limits, whether any worker consumes the task's queue, worker or broker health, and for Kubernetes whether the pod is stuck Pending.

open as a page

What is a zombie task in Airflow, and what causes the scheduler to mark one?

level: seniorimportance: should knowfreq 58%

basics

~20 s

An Airflow zombie task is a task instance the metadata database still shows as running while no live worker process is heartbeating for it. The scheduler detects the missing heartbeats and fails the instance so its retries or downstream handling can proceed.

open as a page

In Airflow, what problem do deferrable operators and the triggerer solve?

level: seniorimportance: should knowfreq 48%

basics

~20 s

They remove the cost of waiting. A deferrable operator starts its work, hands a small async trigger to Airflow's triggerer process and releases its worker slot; the task resumes only when the trigger fires, so thousands of waits cost one process instead of thousands of slots.

open as a page

An Airflow ExternalTaskSensor in a 06:00 DAG times out waiting on a task in a 02:00 DAG. Why?

level: seniorimportance: should knowfreq 40%

basics

~20 s

By default ExternalTaskSensor looks for a run of the other DAG at exactly its own logical date. Two DAGs on different schedules never share a logical date, so the run it waits for does not exist and it waits until timeout.

open as a page

An Airflow DAG paused two weeks with catchup=True is unpaused and floods the scheduler — how do you control it?

level: seniorimportance: should knowfreq 58%

basics

~10 s

Pause the DAG again immediately, then decide which intervals must genuinely be replayed. Cap concurrency with max_active_runs, max_active_tasks_per_dag and a pool, drop the rest, and replay the needed window as an explicit bounded backfill.

open as a page

In Airflow, when does a DAG scheduled on Datasets actually run?

level: seniorimportance: should knowfreq 44%

basics

~20 s

It runs when every Dataset listed in its schedule has been updated at least once since its last run. A Dataset is marked updated when a task that declares it in outlets finishes successfully — Airflow never inspects the underlying storage.

open as a page

showing 1–30 of 37