Data engineering
Part 2 of 6 · Workflow OrchestrationDAG Fundamentals & Airflow Architecture - Topological Scheduling, the Scheduler Loop, Executors, Task States & Pools
Why DAGs and what topological order buys (waves, critical path, cycle detection); Airflow 3.x architecture (Dag processor, Dag bundles, scheduler, API server and Task Execution API, triggerer, metadata DB); task-instance states and trigger rules; pools and concurrency limits; executors compared (Local, Celery, Kubernetes, Edge, multiple executors); HA schedulers via row locks; runnable mini scheduler and wave/mapping simulation; executor decision chart.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
Why must the graph be acyclic?
Answer
A cycle means a task waits on itself, so nothing in it can start.
L2
What does topological order give a scheduler?
Answer
An order plus parallel waves; the critical path sets the minimum run time.
L3
What does the Dag processor do in Airflow 3?
Answer
Parses Dag files from a bundle into serialized Dags so the scheduler never runs author code.
L4
Stuck in scheduled vs stuck in queued?
Answer
Scheduled points at limits like pools and max_active_*; queued points at workers or the executor.
L5
What does a pool protect?
Answer
A shared resource such as a database or API quota, across Dags.
L6
How do multiple schedulers avoid double scheduling?
Answer
Row-level locks in the metadata DB with SKIP LOCKED or NOWAIT.
L7
Celery or Kubernetes executor?
Answer
Celery for many short similar tasks; Kubernetes for isolation and per-task images, at pod start-up cost.
Failure modes
Tasks stuck in queued
The executor accepted them but no worker picked them up: workers down, wrong queue, or pods pending.
Catchup overwhelms a shared database
Many runs start at once without pools or max_active_runs.
Slow parsing from top-level code
Network calls or heavy imports at module level run on every parse and slow every Dag.
Misconceptions
More workers always make a run faster.
Past the critical path, extra slots stop helping.
The scheduler is cheap.
It lives on database round trips, so a slow DB means slow scheduling.
Tasks can still query the metadata DB in Airflow 3.
They talk to the API server through the Task Execution API.
Interviewer traps
Using KubernetesExecutor for thousands of 5-second tasks.
Pod start-up dwarfs the work; batch the tasks or use Celery for that queue.
Touching the scheduler first when tasks sit in queued.
Check workers, queue names and pending pods first.
Design scenario
Same prompt for every reader.
Requirements
Protect the source database and API quota, keep critical Dags fast, survive a scheduler failure, and isolate a few GPU tasks.
Failure assumptions
- A month-long pause is undone with catchup on.
- One scheduler host dies.
- A Dag file has a top-level API call.
Constraints
- PostgreSQL metadata DB.
- Mixed short and heavy tasks.
Prompt
Design the Airflow deployment for 300 Dags that hit a shared Postgres source and a rate-limited SaaS API.
API
Which pools, queues, priority weights and max_active_* settings do Dags declare?
Data
What lives in the metadata DB, and which metrics show scheduler and DB health?
Architecture
Which executors, how many schedulers, and where do the triggerer and Dag processor run?
Overview
A DAG scheduler is a loop over a directed acyclic graph of tasks: find the tasks whose upstream dependencies are satisfied, respect the concurrency limits, hand them to an executor, record what happened, repeat. Airflow is the reference implementation, and its 3.x architecture makes the moving parts explicit: a Dag processor parses your Python files into serialized Dags in a metadata database; the scheduler turns schedules into Dag runs and task instances and moves them through a state machine; an executor (Local, Celery, Kubernetes, Edge, ...) gets them onto workers; a triggerer parks waiting tasks cheaply; and an API server serves the UI and the Task Execution API that workers use instead of touching the database. Once you see that loop, pools, queues, priorities and "why is my task stuck in queued?" stop being mysterious.
Airflow 3.x architecture
| Component | Required? | What it does | What breaks when it is unhealthy |
|---|---|---|---|
| Dag processor | Yes (standalone process since 3.0) | Parses Dag files from a Dag bundle and stores serialized Dags in the metadata DB | New or changed Dags do not appear; import errors pile up |
| Dag bundle | Yes | Where Dag files come from: a local folder by default; GitDagBundle supports versioning so a run can pin the code it started with | Workers may run a different code version than the scheduler planned (non-versioned bundles) |
| Scheduler | Yes | Creates Dag runs from schedules and asset events, evaluates dependencies, queues task instances; the executor runs inside it | Nothing new starts; tasks sit in scheduled |
| Executor | Yes (a scheduler setting) | Gets queued task instances to a place where they run | Tasks sit in queued |
| Workers | Optional for LocalExecutor | Run task code in a supervised subprocess per task instance | Tasks fail or get lost; heartbeats stop |
| Triggerer | Optional | Runs deferred tasks' triggers in an asyncio loop, so waiting costs no worker slot | Deferred tasks never resume |
| API server | Yes | UI and REST API; also the Execution API that the Task SDK uses to report state | UI down; in 3.x tasks cannot report state |
| Metadata DB | Yes | PostgreSQL or MySQL holding Dags, runs, task instances, XComs, variables | Everything stops; it is the single source of truth |
Two 3.x design decisions are worth remembering for interviews. Task code can no longer access the metadata database directly: it talks to the API server through the Task Execution API (AIP-72), which enables remote and multi-language workers (3.3 added experimental Java and Go task SDKs). And the scheduler never runs Dag author code: parsing happens in the Dag processor, so a scheduler is not exposed to a slow or malicious Dag file.
Add worker slots or shorten the critical path?
Prefer
Shorten the critical path
The longest chain sets the floor on run time.
- The critical path list_files > process > join > report took 7 min.
- At 8 slots the run already hit that 7 min floor.
- Splitting or speeding the chain is the only way below it.
Alternative
Keep adding worker slots
More parallelism for the widest wave.
- 1 slot took 26 min, 2 slots 16 min, 4 slots 10 min.
- Gains shrink as waves narrow.
- Past 8 slots more workers stop helping.
The life of a task instance
Diagram 1 condensed: from parse to success, with the failure exits.
- 1
Parse the Dag
The Dag processor stores a serialized, versioned Dag. - 2
Schedule when ready
Upstreams and trigger rules are satisfied; limits are checked. - 3
Queue and run
The executor hands the task to a worker, which reports through the API server. - 4
Retry on failure
up_for_retry waits retry_delay, then the task is scheduled again. - 5
Fail downstream
A hard failure marks dependents upstream_failed; all_done cleanup still runs.
Why a DAG, and what topological order buys you
A DAG says "B needs A's output". Acyclic matters: a cycle means no task can ever start, so every scheduler rejects cycles at parse time. Running tasks in topological order (every task after all of its upstreams) is guaranteed possible only for a DAG, and Kahn's algorithm gives more than an order: it groups tasks into waves that can run in parallel. The critical path (longest chain by duration) is the floor on how long a run can take, however many workers you add.
// Topological scheduling on a static DAG, plus dynamic task mapping.
// 1) Kahn's algorithm groups tasks into "waves" that can run in parallel.
// 2) The critical path (longest path by duration) is the floor on run time,
// no matter how many workers you add.
// 3) With W worker slots, a greedy list scheduler shows the real makespan.
// 4) Dynamic task mapping: the graph SHAPE is fixed at parse time, but one
// node expands into N copies at run time from an upstream result
// (Airflow .expand(), Dagster dynamic outputs, Argo withParam).
type Dag = Record<string, { up: string[]; mins: number }>;
const dag: Dag = {
list_files: { up: [], mins: 1 },
load_dim: { up: [], mins: 4 },
process: { up: ["list_files"], mins: 3 }, // the node we will map
join: { up: ["process", "load_dim"], mins: 2 },
report: { up: ["join"], mins: 1 },
};
function waves(d: Dag): string[][] {
const indeg: Record<string, number> = {}; const users: Record<string, string[]> = {};
for (const n of Object.keys(d)) { indeg[n] = d[n].up.length; users[n] ??= []; for (const u of d[n].up) (users[u] ??= []).push(n); }
let level = Object.keys(d).filter((n) => indeg[n] === 0).sort(); const out: string[][] = []; let seen = 0;
while (level.length) {
out.push(level); seen += level.length; const next: string[] = [];
for (const n of level) for (const v of users[n]) if (--indeg[v] === 0) next.push(v);
level = next.sort();
}
if (seen !== Object.keys(d).length) throw new Error("cycle: not a DAG");
return out;
}
function criticalPath(d: Dag): { mins: number; path: string[] } {
const finish: Record<string, number> = {}; const prev: Record<string, string | null> = {};
for (const n of waves(d).flat()) {
let best = 0, from: string | null = null;
for (const u of d[n].up) if (finish[u] > best) { best = finish[u]; from = u; }
finish[n] = best + d[n].mins; prev[n] = from;
}
let end = Object.keys(finish).reduce((a, b) => (finish[a] >= finish[b] ? a : b));
const path: string[] = []; for (let c: string | null = end; c; c = prev[c]) path.unshift(c);
return { mins: finish[end], path };
}
function makespan(d: Dag, workers: number): number {
// Event-driven greedy scheduler: start any ready task when a slot frees up.
const done = new Map<string, number>(); const running: { n: string; end: number }[] = []; let t = 0;
const pending = new Set(Object.keys(d));
while (pending.size || running.length) {
const ready = [...pending].filter((n) => d[n].up.every((u) => done.has(u) && done.get(u)! <= t)).sort();
while (running.length < workers && ready.length) { const n = ready.shift()!; pending.delete(n); running.push({ n, end: t + d[n].mins }); }
running.sort((a, b) => a.end - b.end); const f = running.shift()!; t = f.end; done.set(f.n, t);
}
return t;
}
function expand(d: Dag, node: string, items: string[], maxMapLength: number): Dag {
if (items.length > maxMapLength) throw new Error(`${node}: ${items.length} mapped tasks > max_map_length ${maxMapLength}, task fails`);
const out: Dag = {};
for (const [n, v] of Object.entries(d)) {
if (n === node) items.forEach((_, i) => (out[`${node}[${i}]`] = { up: v.up, mins: v.mins }));
else out[n] = { up: v.up.flatMap((u) => (u === node ? items.map((_, i) => `${node}[${i}]`) : [u])), mins: v.mins };
}
return out;
}
console.log("parse-time waves:", JSON.stringify(waves(dag)));
const cp = criticalPath(dag); console.log(`critical path: ${cp.path.join(" > ")} = ${cp.mins} min`);
// At run time list_files returns 6 files, so 'process' expands into 6 mapped task instances.
const files = ["a.csv", "b.csv", "c.csv", "d.csv", "e.csv", "f.csv"];
const run = expand(dag, "process", files, 1024);
console.log("run-time waves: ", JSON.stringify(waves(run)));
for (const w of [1, 2, 4, 8]) console.log(` ${w} worker slot(s): makespan ${makespan(run, w)} min`);
console.log(` floor (critical path) ${criticalPath(run).mins} min: more slots stop helping`);
try { expand(dag, "process", files, 4); } catch (e) { console.log("guardrail:", (e as Error).message); }Output:
parse-time waves: [["list_files","load_dim"],["process"],["join"],["report"]]
critical path: list_files > process > join > report = 7 min
run-time waves: [["list_files","load_dim"],["process[0]","process[1]","process[2]","process[3]","process[4]","process[5]"],["join"],["report"]]
1 worker slot(s): makespan 26 min
2 worker slot(s): makespan 16 min
4 worker slot(s): makespan 10 min
8 worker slot(s): makespan 7 min
floor (critical path) 7 min: more slots stop helping
guardrail: process: 6 mapped tasks > max_map_length 4, task failsExpectedparse-time waves: [["list_files","load_dim"],["process"],["join"],["report"]] critical path: list_files > process > join > report = 7 min run-time waves: [["list_files","load_dim"],["process[0]","process[1]","process[2]","process[3]","process[4]","process[5]"],["join"],["report"]] 1 worker slot(s): makespan 26 min 2 worker slot(s): makespan 16 min 4 worker slot(s): makespan 10 min 8 worker slot(s): makespan 7 min floor (critical path) 7 min: more slots stop helping guardrail: process: 6 mapped tasks > max_map_length 4, task fails
Press Run. Snippets must be self-contained — no network, files, or native modules.
Note the split between parse time and run time. The scheduler knows the shape before anything runs, which is what makes backfills and a grid of past runs possible. Dynamic task mapping (Airflow expand(), Dagster dynamic outputs, Argo withParam) keeps that property: the node is declared up front, and only its fan-out count is decided at run time. Airflow caps it with [core] max_map_length, default 1024 mapped tasks; a longer list fails the task.
Diagram 1: the life of a task instance
Decisions
- 1
Step 1: Dag processor parses the Dag file and stores a serialized Dag version
- nextStep 2: scheduler creates a Dag run for the next data interval or asset event
- 2
Step 2: scheduler creates a Dag run for the next data interval or asset event
- nextStep 3: scheduler checks upstream states and the trigger rule
- 3
Step 3: scheduler checks upstream states and the trigger rule
- nextStep 4: task instance moves to scheduled
- nextupstream_failed, task never runs
- 4
Step 4: task instance moves to scheduled
- nextStep 5: free pool slot, Dag and global concurrency limits allow it?
- 5
upstream_failed, task never runs
- ?
Step 5: free pool slot, Dag and global concurrency limits allow it?
- nextStep 4: task instance moves to scheduled
- nextStep 6: executor queues it and a worker starts it, state running
- 7
Step 6: executor queues it and a worker starts it, state running
- nextStep 7: outcome reported through the Execution API
- ?
Step 7: outcome reported through the Execution API
- nextStep 8: success, downstream tasks are re-evaluated
- nextdeferred: trigger waits in the triggerer, no worker slot used
- nextFailure path: up_for_retry, wait retry_delay then back to Step 3
- nextfailed, dependents become upstream_failed unless their trigger rule says otherwise
- 9
Step 8: success, downstream tasks are re-evaluated
- 10
deferred: trigger waits in the triggerer, no worker slot used
- nextStep 4: task instance moves to scheduled
- 11
Failure path: up_for_retry, wait retry_delay then back to Step 3
- 12
failed, dependents become upstream_failed unless their trigger rule says otherwise
Lesson map
DAG Fundamentals & Airflow Architecture - Topological Scheduling, the Scheduler Loop, Executors, Task States & Pools
Why DAGs and what topological order buys (waves, critical path, cycle detection); Airflow 3.x architecture (Dag processor, Dag bundles, scheduler, API server and Task Execution API, triggerer, metadata DB); task-instance states and trigger rules; pools and concurrency limits; executors compared (Local, Celery, Kubernetes, Edge, multiple executors); HA schedulers via row locks; runnable mini scheduler and wave/mapping simulation; executor decision chart.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["Step 1: Dag processor parses the Dag file and stores a serialized Dag version"] b["Step 2: scheduler creates a Dag run for the next data interval or asset event"] c["Step 3: scheduler checks upstream states and the trigger rule"] d["Step 4: task instance moves to scheduled"] uf["upstream_failed, task never runs"] e["Step 5: free pool slot, Dag and global concurrency limits allow it?"] f["Step 6: executor queues it and a worker starts it, state running"] g["Step 7: outcome reported through the Execution API"] h["Step 8: success, downstream tasks are re-evaluated"] t["deferred: trigger waits in the triggerer, no worker slot used"] r["Failure path: up_for_retry, wait retry_delay then back to Step 3"] x["failed, dependents become upstream_failed unless their trigger rule says otherwise"] a -->|continues| b b -->|continues| c c -->|continues| d c -->|continues| uf d -->|continues| e e -->|continues| d e -->|continues| f f -->|continues| g g -->|continues| h g -->|continues| t t -->|continues| d g -->|continues| r g -->|continues| x
Task-instance states
Airflow's states are the vocabulary of every on-call conversation:
| State | Meaning | Usual cause when stuck |
|---|---|---|
none | Not yet queued; dependencies not met | Upstream still running, or Dag run not started |
scheduled | Dependencies met, should run | No pool slot, concurrency limit hit, scheduler busy |
queued | Handed to the executor, waiting for a worker | Workers down or saturated, wrong queue name, pods pending |
running | Executing on a worker | Long task; or a lost worker until heartbeat detection fires |
success / failed | Finished | - |
up_for_retry | Failed, will retry after retry_delay | Working as intended |
up_for_reschedule | A sensor in reschedule mode between pokes | Condition not yet true |
deferred | Waiting on a trigger in the triggerer | Triggerer down, or event never happens |
awaiting_input | Human-in-the-loop task waiting for a person | Nobody has responded |
upstream_failed | An upstream failed and the trigger rule needed it | Fix upstream, then clear |
skipped | Branching, LatestOnly or similar chose not to run it | Working as intended |
removed | Task vanished from the Dag after the run began | Dag code changed mid-run |
The happy path is none -> scheduled -> queued -> running -> success. The runnable scheduler below reproduces it, including pools, a per-Dag concurrency cap, a retry, a hard failure that cascades as upstream_failed, and an all_done cleanup task that runs anyway (Airflow's setup/teardown tasks formalize that pattern).
"""A mini Airflow-style scheduler: DAG parse + cycle check, a scheduler loop,
task-instance states, trigger rules, pools, a concurrency cap and retries.
Real Airflow is far richer (Dag processor, metadata DB, executors, triggerer),
but the state machine below uses Airflow's own state names:
none -> scheduled -> queued -> running -> success | failed | up_for_retry
plus upstream_failed (decided without running) and skipped.
Time is a tick counter so the output is deterministic.
"""
from collections import defaultdict
# ---- the "Dag file": tasks, dependencies and per-task settings ----
DAG = {
# task upstream pool retries delay dur trigger_rule
"extract_orders": dict(up=[], pool="db", retries=0, delay=0, dur=2, rule="all_success"),
"extract_users": dict(up=[], pool="db", retries=0, delay=0, dur=1, rule="all_success"),
"transform": dict(up=["extract_orders", "extract_users"], pool="default", retries=2, delay=2, dur=1, rule="all_success"),
"quality_check": dict(up=["transform"], pool="default", retries=0, delay=0, dur=1, rule="all_success"),
"publish": dict(up=["quality_check"], pool="default", retries=0, delay=0, dur=1, rule="all_success"),
"cleanup_tmp": dict(up=["publish"], pool="default", retries=0, delay=0, dur=1, rule="all_done"),
}
POOLS = {"db": 1, "default": 8} # slots; "db" protects a fragile source database
MAX_ACTIVE_TASKS = 2 # per-Dag concurrency cap
# What happens on each attempt (the outside world): transform fails once, quality_check finds bad data
OUTCOME = {("transform", 1): "fail", ("quality_check", 1): "fail"}
def check_acyclic(dag):
"""DFS three-colour cycle detection, like a Dag processor rejecting a file at parse time."""
color = defaultdict(int)
def visit(n, path):
color[n] = 1
for u in dag[n]["up"]:
if color[u] == 1: raise ValueError("cycle: " + " -> ".join(path + [u]))
if color[u] == 0: visit(u, path + [u])
color[n] = 2
for n in dag: visit(n, [n])
check_acyclic(DAG)
try:
check_acyclic({"a": dict(up=["b"]), "b": dict(up=["a"])})
except ValueError as e:
print("parse error caught:", e)
ti = {t: dict(state="none", attempt=0, ready_at=0, done_at=None) for t in DAG}
TERMINAL = {"success", "failed", "upstream_failed", "skipped"}
def trigger_ok(task):
"""Return 'run', 'wait' or 'upstream_failed' according to the trigger rule."""
ups = [ti[u]["state"] for u in DAG[task]["up"]]
rule = DAG[task]["rule"]
if rule == "all_success":
if any(s in ("failed", "upstream_failed") for s in ups): return "upstream_failed"
return "run" if all(s == "success" for s in ups) else "wait"
if rule == "all_done": # cleanup / teardown style
return "run" if all(s in TERMINAL for s in ups) else "wait"
def set_state(tick, task, new):
old = ti[task]["state"]; ti[task]["state"] = new
print(f"tick {tick:>2} {task:<15} {old:>15} -> {new}")
for tick in range(1, 20):
# 1) finish running tasks whose work is done (worker reports back)
for t, s in ti.items():
if s["state"] == "running" and s["done_at"] == tick:
if OUTCOME.get((t, s["attempt"])) == "fail":
if s["attempt"] <= DAG[t]["retries"]:
s["ready_at"] = tick + DAG[t]["delay"]; set_state(tick, t, "up_for_retry")
else:
set_state(tick, t, "failed")
else:
set_state(tick, t, "success")
# 2) queued -> running (a worker picked the task up)
for t, s in ti.items():
if s["state"] == "queued":
s["attempt"] += 1; s["done_at"] = tick + DAG[t]["dur"]; set_state(tick, t, "running")
# 3) scheduler: evaluate dependencies for none / up_for_retry tasks
for t, s in ti.items():
if s["state"] in ("none", "up_for_retry") and tick >= s["ready_at"]:
verdict = trigger_ok(t)
if verdict == "run": set_state(tick, t, "scheduled")
elif verdict == "upstream_failed": set_state(tick, t, "upstream_failed")
# 4) executor: scheduled -> queued if a pool slot and the Dag concurrency cap allow it
active = sum(s["state"] in ("queued", "running") for s in ti.values())
used = defaultdict(int)
for t, s in ti.items():
if s["state"] in ("queued", "running"): used[DAG[t]["pool"]] += 1
for t, s in ti.items():
if s["state"] == "scheduled":
p = DAG[t]["pool"]
if used[p] < POOLS[p] and active < MAX_ACTIVE_TASKS:
used[p] += 1; active += 1; set_state(tick, t, "queued")
if all(s["state"] in TERMINAL for s in ti.values()):
break
failed = any(s["state"] in ("failed", "upstream_failed") for s in ti.values())
print("\nDag run state:", "failed" if failed else "success")
print({t: s["state"] for t, s in ti.items()})Output:
parse error caught: cycle: a -> b -> a
tick 1 extract_orders none -> scheduled
tick 1 extract_users none -> scheduled
tick 1 extract_orders scheduled -> queued
tick 2 extract_orders queued -> running
tick 4 extract_orders running -> success
tick 4 extract_users scheduled -> queued
tick 5 extract_users queued -> running
tick 6 extract_users running -> success
tick 6 transform none -> scheduled
tick 6 transform scheduled -> queued
tick 7 transform queued -> running
tick 8 transform running -> up_for_retry
tick 10 transform up_for_retry -> scheduled
tick 10 transform scheduled -> queued
tick 11 transform queued -> running
tick 12 transform running -> success
tick 12 quality_check none -> scheduled
tick 12 quality_check scheduled -> queued
tick 13 quality_check queued -> running
tick 14 quality_check running -> failed
tick 14 publish none -> upstream_failed
tick 14 cleanup_tmp none -> scheduled
tick 14 cleanup_tmp scheduled -> queued
tick 15 cleanup_tmp queued -> running
tick 16 cleanup_tmp running -> success
Dag run state: failed
{'extract_orders': 'success', 'extract_users': 'success', 'transform': 'success', 'quality_check': 'failed', 'publish': 'upstream_failed', 'cleanup_tmp': 'success'}Concurrency controls: who limits what
| Knob | Scope | Use it for |
|---|---|---|
[core] parallelism | Whole installation, per scheduler | Upper bound on running task instances |
max_active_tasks (Dag) | One Dag across its runs | Keep one Dag from hogging workers |
max_active_runs (Dag) | Runs of one Dag | Limit how many intervals run at once, especially during catchup |
max_active_tis_per_dag (task) | One task across runs | A task that must not overlap with itself |
| Pools | Any set of tasks, across Dags | Protect a shared resource: a database, an API quota. default_pool starts with 128 slots |
pool_slots (task) | One task instance | Heavy tasks that should count as several slots |
queue (task) | Celery queues and similar | Route GPU or high-memory tasks to the right workers |
priority_weight and weight_rule | Ordering when slots are scarce | Critical paths first; in 3.x priority is capped by pool slots |
Executors compared
| Executor | Where tasks run | Strengths | Weaknesses | Pick when |
|---|---|---|---|---|
| LocalExecutor | Subprocesses of the scheduler host | Simple, low latency, no broker | Shares one machine with the scheduler | Small installs, dev, CI (SequentialExecutor was removed in 3.0) |
| CeleryExecutor | Long-running worker processes fed by a broker (Redis or RabbitMQ) | Low start latency, many small tasks, autoscalable worker fleet | Broker to run, noisy neighbours, shared Python environment | Many short tasks with similar dependencies |
| KubernetesExecutor | One pod per task instance | Isolation, per-task images and resources, scales to zero | Pod start-up latency per task, cluster to operate | Heterogeneous tasks, strict isolation |
| EdgeExecutor | Workers in remote or edge sites over HTTP | Run close to data or behind firewalls | Newer, more moving parts | Hybrid and remote compute |
| Batch / ECS-style containerized executors | Cloud batch services | Managed capacity | Cloud-specific, start-up latency | Already on that cloud service |
| Multiple executors | Mix per task (executor= on a task) | Fast tasks on Celery, heavy ones on Kubernetes | More configuration | Mixed workloads |
The scheduler loop and HA
The scheduler runs a loop: create due Dag runs, examine runnable task instances, check limits, queue them through the executor, and process executor events (finished, failed, lost). Several schedulers can run at once for throughput and resilience; they coordinate through the metadata database with row-level locks (SELECT ... FOR UPDATE with SKIP LOCKED or NOWAIT), which is why PostgreSQL 12+ or MySQL 8+ is required for multiple schedulers. That design choice means the database is both the coordination service and the bottleneck: watch its CPU, connection count and lock waits.
Dag file parsing is a separate, continuous cost. The Dag processor re-parses files on an interval (min_file_process_interval), so any work at module top level (network calls, database queries, heavy imports) runs again and again, slowing parsing for every Dag in the bundle. Airflow's own best-practices guide is blunt: no database access, heavy computation or networking at top level.
Illustration: an Airflow 3 Dag (not executed here)
from datetime import timedelta
import pendulum
from airflow.sdk import dag, task
@dag(
schedule="@daily",
start_date=pendulum.datetime(2026, 1, 1, tz="UTC"),
catchup=False, # 3.x default; set it explicitly anyway
max_active_runs=2,
default_args={"retries": 3, "retry_delay": timedelta(minutes=5)},
)
def orders_daily():
@task(pool="orders_db") # pool protects the source database
def list_files(**context) -> list[str]:
start = context["data_interval_start"]
return [f"s3://raw/orders/{start:%Y-%m-%d}/part-{i}.parquet" for i in range(4)]
@task
def process(path: str) -> str:
return path.replace("raw", "clean") # pass paths, not data
@task
def publish(paths: list[str]) -> None: ...
publish(process.expand(path=list_files())) # dynamic task mapping
orders_daily()Diagram 2: picking an executor
Decisions
- 1
Step 1: estimate task volume, isolation needs and where workers must run
- nextStep 2: one machine is enough?
- ?
Step 2: one machine is enough?
- nextLocalExecutor
- nextStep 3: tasks need per-task images, resources or strong isolation?
- 3
LocalExecutor
- nextFailure path: outgrowing one host shows as tasks stuck in queued, move to Step 3
- ?
Step 3: tasks need per-task images, resources or strong isolation?
- nextKubernetesExecutor, one pod per task
- nextStep 4: many short tasks, low start-up latency matters?
- 5
KubernetesExecutor, one pod per task
- nextStep 6: some Dags need a different trade-off?
- ?
Step 4: many short tasks, low start-up latency matters?
- nextCeleryExecutor with a broker and a warm worker fleet
- nextStep 5: workers in remote sites or other networks?
- 7
CeleryExecutor with a broker and a warm worker fleet
- nextStep 6: some Dags need a different trade-off?
- ?
Step 5: workers in remote sites or other networks?
- nextEdgeExecutor, workers pull over HTTP
- nextCeleryExecutor with a broker and a warm worker fleet
- 9
EdgeExecutor, workers pull over HTTP
- ?
Step 6: some Dags need a different trade-off?
- nextConfigure multiple executors and set executor per task or Dag
- nextDone
- 11
Configure multiple executors and set executor per task or Dag
- 12
Done
- 13
Failure path: outgrowing one host shows as tasks stuck in queued, move to Step 3
What happens if you choose otherwise
- KubernetesExecutor for thousands of 5-second tasks: pod scheduling and image pull time can dwarf the work. Batch the work into fewer tasks or use Celery for that queue.
- CeleryExecutor for tasks with conflicting dependencies: every worker shares one Python environment; you end up with dependency hell or per-queue images. Use KubernetesExecutor,
@task.virtualenv, or container operators. - No pools around a shared database: catchup or a backfill launches many runs at once and the source falls over. Pools and
max_active_runsare the circuit breaker. - One giant Dag with thousands of tasks: parsing, the grid view and scheduling all slow down. Split by domain and connect with assets.
Pitfalls
- Tasks stuck in
queued: the executor accepted them but no worker took them. Check worker health, queue names, and Kubernetes pending pods before touching the scheduler. - Tasks stuck in
scheduled: limits, not workers. Check pools,max_active_tasks,max_active_runsandparallelism. - Changing a Dag mid-run: tasks can become
removed, and on a non-versioned bundle workers may run newer code than the run started with. Prefer a versioned bundle such asGitDagBundle. - Assuming the scheduler is cheap: it lives on database round trips. Slow DB equals slow scheduling.
Interview Q&A
Walk through how a task gets from a Python file to running code in Airflow 3.
Answer
The Dag processor parses the file from a Dag bundle and stores a serialized, versioned Dag in the metadata DB. The scheduler creates a Dag run for the due interval, evaluates each task's upstream states and trigger rule, moves ready ones to scheduled, then to queued when limits allow and hands them to the executor. A worker runs the task in a supervised subprocess and reports state through the API server's Execution API.
Why must the graph be acyclic, and how is that checked?
Answer
A cycle means some task waits on itself, so nothing in the cycle can start. Topological sort (Kahn's algorithm or DFS with colouring) detects it at parse time, and the Dag fails to import.
A Dag's tasks sit in `scheduled` for an hour. What do you check?
Answer
Concurrency limits first: pool slots, the Dag's max_active_tasks and max_active_runs, global parallelism, and whether a catchup or backfill is consuming them. Then scheduler health and database latency.
How do multiple Airflow schedulers avoid scheduling the same task twice?
Answer
They coordinate through the metadata database with row-level locks (FOR UPDATE ... SKIP LOCKED or NOWAIT), so each scheduler claims distinct Dag runs and task instances. That is why multiple schedulers need PostgreSQL 12+ or MySQL 8+.
What does the triggerer buy you?
Answer
A deferred task releases its worker slot and its waiting moves to an async trigger in the triggerer process, where thousands of waits share one event loop. Waiting stops costing a worker per wait.
What does max_map_length guard against?
Answer
Dynamic task mapping producing too many task instances; by default more than 1024 mapped tasks fails the task.
Why does Airflow 3 move parsing out of the scheduler?
Answer
So the scheduler never runs Dag author code; a slow or malicious Dag file cannot stall scheduling.
What is the all_done trigger rule for?
Answer
Running a task such as cleanup whatever its upstreams ended as, which is how cleanup_tmp still ran after quality_check failed.
Check yourself
Take a Dag you know, draw its waves with Kahn's algorithm, estimate each task's duration, and compute the critical path and the slot count beyond which more workers stop helping.
Elsewhere in the library
These pages stay as they are. This lesson only points at them: Graphs — Topological Sort & DAGs — Kahn, DFS Finish Times & Cycles, Dependency Graphs & Incremental Builds - Task DAGs, Topological Scheduling, Content Hashing & Affected Detection, Row Locks: SELECT FOR UPDATE, Bulkheads — Thread Pools, Connection Pools, Queues & Failure Domains, Kubernetes Workloads — Deployments, Probes, Resources & Progressive Delivery, Distributed Locks — Correctness, Leases & Fencing Tokens.