Data engineering
Part 1 of 6 · Workflow OrchestrationWorkflow Orchestration - Scheduled DAGs vs Durable Execution (Airflow, Temporal & Friends)
Concept hub: two families of orchestrators, scheduled batch DAGs (Airflow, Dagster, Prefect, Argo) vs durable execution (Temporal, Cadence, Step Functions, Durable Functions, Restate, Inngest); the shared kernel (durable state, scheduler, queue, workers, retries, heartbeats/leases, at-least-once units so idempotency matters); comparison tables; runnable task-level vs step-journal crash demo and lease/heartbeat kernel; decision chart.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What do all orchestrators have in common?
Answer
Durable state, a scheduler, a queue, workers, retry policies and heartbeat or timeout liveness.
L2
Why is every unit of work at-least-once?
Answer
The orchestrator can only guess a worker died. The 'dead' worker may still finish, so work can run twice.
L3
What does a run represent in each family?
Answer
A data interval or partition for DAG tools; a business entity or event for durable engines.
L4
How is the graph known?
Answer
Declared and parsed up front for DAG tools; discovered as code runs for durable engines.
L5
What happens after a crash mid-run?
Answer
A DAG tool reruns the whole task; a durable engine replays the code and skips recorded steps.
L6
When would you choose Step Functions over Temporal?
Answer
All-in on AWS, mostly AWS service calls, no servers to run, and logic that fits a state machine.
L7
Is it normal to run both families in one company?
Answer
Yes: Airflow or Dagster for data, Temporal or Step Functions for product workflows.
Failure modes
Duplicate side effects after a presumed-dead worker
The orchestrator retries while the original worker still finishes, so non-idempotent work runs twice.
Request-path workflows on a batch scheduler
Each run pays scheduler and database overhead and has no signals or step journal.
The orchestrator used as a data plane
Large payloads in XCom or workflow history bloat the state store everything depends on.
Misconceptions
Orchestrators give exactly-once side effects.
They record each step once in their own state. The effect inside a step runs at least once.
One family is simply better.
They answer different questions: interval batch vs crash-proof business processes.
A durable workflow is the same as a saga.
A saga is a pattern; a durable engine is a good place to run an orchestrated saga.
Interviewer traps
Picking an orchestrator by popularity.
Write the one-sentence description of a run first, then match the family.
Passing DataFrames between steps.
Pass object-store paths or references and keep payloads small.
Design scenario
Same prompt for every reader.
Requirements
Backfillable nightly loads by date, crash-proof onboarding with human approvals and timeouts, and one place to see where each item is stuck.
Failure assumptions
- Workers are killed during deploys.
- An approval arrives while no worker is running.
- A nightly run is missed for a week.
Constraints
- Payloads must stay small.
- Side effects hit external APIs without built-in deduplication.
Prompt
A company runs nightly warehouse loads and also needs a multi-day customer onboarding flow with approvals. Design the orchestration layer.
API
How do product services start or signal onboarding runs, and how do data tasks get their intervals?
Data
Where does each family keep state, and which idempotency keys protect side effects?
Architecture
Which orchestrator owns each job, and how do the two hand off work?
Overview
"Orchestrator" covers two different machines that share one kernel. A scheduled batch DAG system (Airflow, Dagster, Prefect, Argo Workflows) answers "for each time window or data partition, which tasks must run, in which order, and did they all succeed?" A durable execution engine (Temporal, Cadence, AWS Step Functions, Azure Durable Functions, Restate, Inngest) answers "this business process must run to completion even if every process running it crashes, possibly over days." Both keep state outside the worker, retry failed work, run a scheduler that hands tasks to workers through a queue, and detect dead workers with heartbeats or timeouts. They differ in what a run is (a data interval vs a business event), how the graph is known (declared up front vs discovered as code runs), and how fine the saved state is (per task vs per step). Learn those differences and you can place any new tool in one sentence.
The shared kernel
Strip the branding and every orchestrator has the same five parts:
| Part | Job | Airflow name | Temporal name | Step Functions name |
|---|---|---|---|---|
| Durable state store | Survives every crash; source of truth for progress | Metadata database (PostgreSQL or MySQL) | Persistence store (Cassandra, PostgreSQL, MySQL) holding mutable state and event history | Managed by AWS (execution history) |
| Scheduler | Decides what is runnable now | Scheduler loop (plus Dag processor that parses files) | History service decides; Matching service dispatches | The state machine interpreter |
| Queue | Decouples deciding from doing | Executor queue (Celery broker, Kubernetes pods, ...) | Task Queues | Internal, plus activity task polling |
| Workers | Run your code | Celery workers, Kubernetes pods, local processes | Your worker processes polling a Task Queue | Lambda, ECS, any integrated service, or activity workers |
| Retry and liveness | Retry failures, notice silent deaths | retries, retry_delay, task heartbeats | Retry Policy, Start-To-Close and Heartbeat timeouts | Retry and Catch fields, TimeoutSeconds, HeartbeatSeconds |
Recovering a run that crashed mid-way
Prefer
Step-level journal (durable execution)
Record every step's result and replay the code after a crash.
- order-42 resumed after charge_card.
- Every side effect ran exactly once in the demo.
- The cost: workflow code must be deterministic.
Alternative
Task-level state (scheduled DAG)
Record each task's state and rerun a failed task from its start.
- transform reran from its first step, so clean ran twice.
- Whole-partition reruns and backfills stay simple and visible.
- Work inside a task must be safe to repeat.
One run through either family
Diagram 1 condensed: the shared kernel and where the families split.
- 1
Trigger a run
A schedule or data event for DAGs; a business event for durable engines. - 2
Lease work to a worker
The scheduler queues runnable units and a worker picks them up. - 3
Detect silent death
A missed heartbeat or timeout triggers a retry with backoff. - 4
Recover
Rerun the task, or replay the code and skip recorded steps. - 5
Repeat a side effect
The presumed-dead worker may still finish, so units must be idempotent.
The runnable model below shows the kernel without any branding: a scheduler leases tasks to workers, a retry policy backs off exponentially, and a heartbeat timeout notices a worker that died without saying anything.
// The kernel every orchestrator shares, whichever family it belongs to:
// durable state store + scheduler loop + task queue + workers + retry policy
// + a heartbeat/lease so a silently dead worker is detected.
// Airflow calls the dead case a task whose heartbeat stopped (zombie/orphan task);
// Temporal calls it a Heartbeat or Start-To-Close timeout. Same idea.
// Simulated clock in seconds, so output is deterministic.
type State = "scheduled" | "running" | "retry_wait" | "success" | "failed";
interface Task { id: string; state: State; attempt: number; lastBeat: number; notBefore: number; worker?: string }
interface Policy { maxAttempts: number; initial: number; coeff: number; maxInterval: number; heartbeatTimeout: number }
const policy: Policy = { maxAttempts: 4, initial: 1, coeff: 2, maxInterval: 10, heartbeatTimeout: 5 };
const backoff = (attempt: number) => Math.min(policy.initial * policy.coeff ** (attempt - 1), policy.maxInterval);
const tasks: Task[] = [
{ id: "resize-images", state: "scheduled", attempt: 0, lastBeat: 0, notBefore: 0 },
{ id: "charge-card", state: "scheduled", attempt: 0, lastBeat: 0, notBefore: 0 },
];
// Script of what each worker does on a given attempt (the "real world").
const behaviour: Record<string, (attempt: number) => "ok" | "error" | "dies"> = {
"resize-images": (a) => (a === 1 ? "dies" : "ok"), // worker host vanishes on attempt 1
"charge-card": (a) => (a < 3 ? "error" : "ok"), // flaky dependency, succeeds on attempt 3
};
const workDuration = 3; // seconds each attempt takes
const inflight = new Map<string, { finishAt: number; outcome: "ok" | "error" | "dies" }>();
const log = (t: number, msg: string) => console.log(`t=${String(t).padStart(2)}s ${msg}`);
for (let now = 0; now <= 30; now++) {
// 1. Scheduler loop: dispatch runnable tasks to the queue -> a worker picks it up.
for (const task of tasks) {
if ((task.state === "scheduled" || task.state === "retry_wait") && now >= task.notBefore) {
task.attempt++; task.state = "running"; task.lastBeat = now; task.worker = `${task.id.split("-")[0]}-w${task.attempt}`;
const outcome = behaviour[task.id](task.attempt);
inflight.set(task.id, { finishAt: now + workDuration, outcome });
log(now, `${task.id}: attempt ${task.attempt} leased to ${task.worker}`);
}
}
// 2. Workers heartbeat while alive, and report results when done.
for (const task of tasks.filter((t) => t.state === "running")) {
const job = inflight.get(task.id)!;
if (job.outcome !== "dies") task.lastBeat = now;
if (now === job.finishAt && job.outcome !== "dies") {
inflight.delete(task.id);
if (job.outcome === "ok") { task.state = "success"; log(now, `${task.id}: success`); }
else fail(task, now, "activity error");
}
}
// 3. Lease check: no heartbeat within the timeout means the worker is presumed dead.
for (const task of tasks.filter((t) => t.state === "running")) {
if (now - task.lastBeat > policy.heartbeatTimeout) {
inflight.delete(task.id);
fail(task, now, `no heartbeat for ${now - task.lastBeat}s (worker ${task.worker} presumed dead)`);
}
}
}
function fail(task: Task, now: number, why: string) {
if (task.attempt >= policy.maxAttempts) { task.state = "failed"; log(now, `${task.id}: FAILED for good (${why})`); return; }
const wait = backoff(task.attempt);
task.state = "retry_wait"; task.notBefore = now + wait;
log(now, `${task.id}: ${why} -> retry in ${wait}s`);
}
console.log("final:", tasks.map((t) => `${t.id}=${t.state} after ${t.attempt} attempts`).join(", "));
console.log("note: the dead worker may still finish resize-images later -> tasks must be idempotent");Output:
t= 0s resize-images: attempt 1 leased to resize-w1
t= 0s charge-card: attempt 1 leased to charge-w1
t= 3s charge-card: activity error -> retry in 1s
t= 4s charge-card: attempt 2 leased to charge-w2
t= 6s resize-images: no heartbeat for 6s (worker resize-w1 presumed dead) -> retry in 1s
t= 7s resize-images: attempt 2 leased to resize-w2
t= 7s charge-card: activity error -> retry in 2s
t= 9s charge-card: attempt 3 leased to charge-w3
t=10s resize-images: success
t=12s charge-card: success
final: resize-images=success after 2 attempts, charge-card=success after 3 attempts
note: the dead worker may still finish resize-images later -> tasks must be idempotentExpectedt= 0s resize-images: attempt 1 leased to resize-w1 t= 0s charge-card: attempt 1 leased to charge-w1 t= 3s charge-card: activity error -> retry in 1s t= 4s charge-card: attempt 2 leased to charge-w2 t= 6s resize-images: no heartbeat for 6s (worker resize-w1 presumed dead) -> retry in 1s t= 7s resize-images: attempt 2 leased to resize-w2 t= 7s charge-card: activity error -> retry in 2s t= 9s charge-card: attempt 3 leased to charge-w3 t=10s resize-images: success t=12s charge-card: success final: resize-images=success after 2 attempts, charge-card=success after 3 attempts note: the dead worker may still finish resize-images later -> tasks must be idempotent
Press Run. Snippets must be self-contained — no network, files, or native modules.
The last line is the most important sentence in this whole series: an orchestrator can only guess that a worker died. The "dead" worker may still be running. So every orchestrator gives you at-least-once execution of the unit of work, and correctness comes from making that unit idempotent.
Where the two families differ
| Dimension | Scheduled batch DAGs | Durable execution |
|---|---|---|
| What starts a run | A schedule over time (a data interval) or a data event (an asset or partition update) | A business event: an API call, a message, a signal, a cron schedule |
| What a run represents | "Process the data for 2026-10-06" | "Fulfil order 42", "onboard user 9", "run this agent loop" |
| How the graph is known | Declared up front in a Dag file and parsed before running; dynamic mapping can fan out at run time | Discovered as ordinary code executes: loops, ifs and function calls are the graph |
| Unit of saved state | A task instance (success, failed, up_for_retry, ...) | Every step: each activity result, timer, signal and child workflow is an event in history |
| Recovery after a crash | Retry the whole task from its start | Replay the code; finished steps return their recorded result |
| Typical duration | Minutes to hours per run | Milliseconds to months (durable timers, human waits) |
| Typical latency to start | Seconds or more; tuned for throughput | Low, built for request paths and event handling |
| Data between steps | Small metadata via XCom-style channels; big data goes to storage | Arguments and results are persisted in history, with size limits (Temporal errors at 2 MB per payload, Step Functions caps state input and output at 256 KiB) |
| Who it is written for | Data engineers, analytics, ML pipelines | Application and platform engineers, backend services |
| Mental model | Make or cron for data, with dependencies and backfills | A function call that cannot be killed |
The same three-step job behaves very differently in each family when a crash lands mid-way:
"""The same 3-step job run by the two orchestration families.
Family A, scheduled batch DAG (Airflow / Dagster / Prefect / Argo style):
- a run exists per data interval; state is tracked per TASK
- a crash inside a task means the WHOLE task is retried from its start
Family B, durable execution (Temporal / Step Functions / Restate / Inngest style):
- a run exists per business event; state is a journal of completed STEPS
- after a crash the function is re-run, but completed steps return their
recorded result instead of executing again
Everything is simulated in-process so the output is deterministic.
"""
calls = {} # side-effect counter: how many times each real action executed
def side_effect(name):
calls[name] = calls.get(name, 0) + 1
class Crash(Exception):
pass
# ---------------- Family A: task-level state ----------------
def run_dag(interval, crash_in_task=None, crash_after_step=None, state=None):
"""Tasks: extract -> transform -> load. Each task is 2 internal steps."""
state = state if state is not None else {}
tasks = {
"extract": ["download_part_1", "download_part_2"],
"transform": ["clean", "aggregate"],
"load": ["write_partition", "publish_metrics"],
}
for task, steps in tasks.items(): # topological order
if state.get(task) == "success":
print(f" [{interval}] {task:<9} skipped (already success)")
continue
try:
for step in steps: # the orchestrator cannot see these
side_effect(step)
if task == crash_in_task and step == crash_after_step:
raise Crash(step)
state[task] = "success"
print(f" [{interval}] {task:<9} success")
except Crash as c:
state[task] = "up_for_retry"
print(f" [{interval}] {task:<9} CRASH after {c} -> up_for_retry (whole task reruns)")
return state
return state
# ---------------- Family B: step-level journal ----------------
def durable_order_workflow(ctx, order_id):
"""Ordinary-looking code; every ctx.step is journaled by the engine."""
ctx.step("reserve_stock", order_id)
ctx.step("charge_card", order_id)
ctx.step("create_shipment", order_id)
ctx.step("send_email", order_id)
return "done"
class DurableCtx:
def __init__(self, journal, crash_after=None):
self.journal, self.crash_after, self.i = journal, crash_after, 0
def step(self, name, arg):
if self.i < len(self.journal): # replay: return recorded result
rec = self.journal[self.i]; self.i += 1
assert rec[0] == name, "non-deterministic workflow code"
return rec[1]
side_effect(name) # first time: really execute
result = f"{name}:{arg}:ok"
self.journal.append((name, result)); self.i += 1
if name == self.crash_after:
raise Crash(name)
return result
def run_durable(order_id, crash_after=None, journal=None):
journal = journal if journal is not None else []
try:
out = durable_order_workflow(DurableCtx(journal, crash_after), order_id)
print(f" [{order_id}] completed, journal={[j[0] for j in journal]}")
return journal, out
except Crash as c:
print(f" [{order_id}] CRASH after step {c}; journal has {len(journal)} steps")
return journal, None
print("Family A: scheduled DAG, task-level state")
st = run_dag("2026-10-06", crash_in_task="transform", crash_after_step="aggregate")
st = run_dag("2026-10-06", state=st) # retry of the same data interval
print(" side effects:", {k: calls[k] for k in ["download_part_1", "clean", "aggregate", "write_partition"]})
calls.clear()
print("\nFamily B: durable execution, step-level journal")
j, _ = run_durable("order-42", crash_after="charge_card")
j, _ = run_durable("order-42", journal=j) # worker restarts, replays journal
print(" side effects:", calls)
print("\nTakeaway: A retried 'transform' from its first step (clean ran twice);")
print("B resumed after charge_card without charging twice.")Output:
Family A: scheduled DAG, task-level state
[2026-10-06] extract success
[2026-10-06] transform CRASH after aggregate -> up_for_retry (whole task reruns)
[2026-10-06] extract skipped (already success)
[2026-10-06] transform success
[2026-10-06] load success
side effects: {'download_part_1': 1, 'clean': 2, 'aggregate': 2, 'write_partition': 1}
Family B: durable execution, step-level journal
[order-42] CRASH after step charge_card; journal has 2 steps
[order-42] completed, journal=['reserve_stock', 'charge_card', 'create_shipment', 'send_email']
side effects: {'reserve_stock': 1, 'charge_card': 1, 'create_shipment': 1, 'send_email': 1}
Takeaway: A retried 'transform' from its first step (clean ran twice);
B resumed after charge_card without charging twice.Neither is "better". The DAG system's coarse checkpoints make reruns and backfills of whole partitions trivial and visible. The durable engine's fine journal lets a payment workflow continue after charge_card without charging twice.
Diagram 1: one run through either family
Decisions
- 1
Step 1: trigger arrives (schedule tick, data event, API call or signal)
- nextStep 2: orchestrator writes a new run to its durable state store
- 2
Step 2: orchestrator writes a new run to its durable state store
- nextStep 3: scheduler picks the next runnable unit (task or step)
- 3
Step 3: scheduler picks the next runnable unit (task or step)
- nextStep 4: unit goes onto a queue and a worker leases it
- 4
Step 4: unit goes onto a queue and a worker leases it
- nextStep 5: did the worker report success before its timeout?
- ?
Step 5: did the worker report success before its timeout?
- nextStep 6: record the result (task state or history event)
- nextFailure path: retry policy decides retry with backoff or give up
- 6
Step 6: record the result (task state or history event)
- nextStep 7: more units left?
- ?
Step 7: more units left?
- nextStep 3: scheduler picks the next runnable unit (task or step)
- nextStep 8: run is complete, downstream runs or callers are notified
- 8
Step 8: run is complete, downstream runs or callers are notified
- 9
Failure path: retry policy decides retry with backoff or give up
- nextStep 4: unit goes onto a queue and a worker leases it
- nextRun fails: alert, compensate, or wait for a human to clear and rerun
- 10
Run fails: alert, compensate, or wait for a human to clear and rerun
Lesson map
Workflow Orchestration - Scheduled DAGs vs Durable Execution (Airflow, Temporal & Friends)
Concept hub: two families of orchestrators, scheduled batch DAGs (Airflow, Dagster, Prefect, Argo) vs durable execution (Temporal, Cadence, Step Functions, Durable Functions, Restate, Inngest); the shared kernel (durable state, scheduler, queue, workers, retries, heartbeats/leases, at-least-once units so idempotency matters); comparison tables; runnable task-level vs step-journal crash demo and lease/heartbeat kernel; 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: trigger arrives (schedule tick, data event, API call or signal)"] b["Step 2: orchestrator writes a new run to its durable state store"] c["Step 3: scheduler picks the next runnable unit (task or step)"] d["Step 4: unit goes onto a queue and a worker leases it"] e["Step 5: did the worker report success before its timeout?"] f["Step 6: record the result (task state or history event)"] g["Step 7: more units left?"] h["Step 8: run is complete, downstream runs or callers are notified"] x["Failure path: retry policy decides retry with backoff or give up"] y["Run fails: alert, compensate, or wait for a human to clear and rerun"] a -->|continues| b b -->|continues| c c -->|continues| d d -->|continues| e e -->|continues| f f -->|continues| g g -->|continues| c g -->|continues| h e -->|continues| x x -->|continues| d x -->|continues| y
The tools on one page
| Tool | Family | Graph style | State granularity | Runs on | Note |
|---|---|---|---|---|---|
| Apache Airflow 3.x | Scheduled DAG | Python Dag files parsed by a Dag processor; dynamic task mapping | Task instance | Your infra or managed (MWAA, Cloud Composer, Astronomer) | 3.0 added Dag versioning, assets, a Task Execution API; 3.2 and 3.3 added asset partitions |
| Dagster | Scheduled DAG, asset-first | Software-defined assets and their dependencies | Asset materialization per partition | Your infra or Dagster+ | You declare what data should exist; runs are a means |
| Prefect 3 | Scheduled DAG, code-first | Plain Python flows and tasks; graph discovered at run time | Task run, with caching | Your infra via work pools, or Prefect Cloud | Feels closest to durable code among the DAG tools |
| Argo Workflows | Scheduled DAG on Kubernetes | YAML dag or steps templates, each step a pod | Node (pod) | Kubernetes | Great for container-per-step batch and ML |
| Temporal | Durable execution | Workflow code in Go, Java, Python, TypeScript, .NET and more | Event per step | Self-hosted cluster or Temporal Cloud | Event history plus deterministic replay |
| Cadence | Durable execution | Same model as Temporal (Temporal forked from it) | Event per step | Self-hosted | Originated at Uber |
| AWS Step Functions | Durable execution, declarative | JSON/YAML state machine (Amazon States Language) | State transition | Managed by AWS | Standard: up to one year, exactly-once workflow execution; Express: up to five minutes, at-least-once (async) |
| Azure Durable Functions | Durable execution | Orchestrator functions in code | Event per step (event sourcing) | Azure Functions | Orchestrators replay, so they must be deterministic |
| Restate | Durable execution | Handlers in ordinary services | Journal entry per step | Restate server in front of your services | Single Rust binary; services stay normal HTTP apps |
| Inngest | Durable execution | step.run blocks inside functions | Checkpoint per step | Inngest platform calling your endpoints | Waits do not hold compute |
Diagram 2: decision chart
Decisions
- 1
Step 1: describe one run in a sentence
- nextStep 2: is a run processing data for a time window or partition?
- ?
Step 2: is a run processing data for a time window or partition?
- nextStep 3: do you think in tables and datasets more than in tasks?
- nextStep 5: must a business process survive crashes, wait on humans or timers, or compensate?
- ?
Step 3: do you think in tables and datasets more than in tasks?
- nextAsset-oriented DAG tool (Dagster, Airflow assets)
- nextStep 4: every step is a container on Kubernetes?
- 4
Asset-oriented DAG tool (Dagster, Airflow assets)
- ?
Step 4: every step is a container on Kubernetes?
- nextArgo Workflows
- nextTask DAG scheduler (Airflow, Prefect)
- 6
Argo Workflows
- 7
Task DAG scheduler (Airflow, Prefect)
- ?
Step 5: must a business process survive crashes, wait on humans or timers, or compensate?
- nextStep 6: all-AWS, mostly service calls, happy with a JSON state machine?
- nextPlain cron plus a queue with idempotent consumers is enough
- ?
Step 6: all-AWS, mostly service calls, happy with a JSON state machine?
- nextStep Functions
- nextCode-first durable execution (Temporal, Restate, Inngest, Durable Functions)
- 10
Step Functions
- 11
Code-first durable execution (Temporal, Restate, Inngest, Durable Functions)
- 12
Plain cron plus a queue with idempotent consumers is enough
- nextFailure path: once you add timers, compensations and per-item status by hand, go back to Step 5
- 13
Failure path: once you add timers, compensations and per-item status by hand, go back to Step 5
What happens if you choose otherwise
- Airflow as a request-path workflow engine (one Dag run per order): every run carries scheduler, executor and database overhead tuned for batch throughput, and the model has no first-class signals, queries or per-step journal. Teams end up building a state machine in XCom and sensors.
- Temporal for nightly warehouse loads: it works, but you rebuild what a DAG scheduler gives you for free: data intervals, catchup and backfill over date ranges, a grid of partitions in the UI, and data lineage. You also must keep big data out of history.
- Step Functions for logic-heavy flows: branching and loops become large JSON documents, state transitions are billed on Standard workflows, and Standard executions fail at 25,000 history events unless you split them.
- Cron plus a queue for a multi-step, multi-day process: possible, but you now own retries, timeouts, compensation, visibility, and "where is order 42 stuck?". That is exactly what durable execution sells.
- Two orchestrators in one company is normal: Airflow or Dagster for data, Temporal or Step Functions for product workflows. A DAG task can start a workflow, and a workflow activity can trigger a Dag run.
Pitfalls
- Believing "exactly-once". Orchestrators record a step exactly once in their own state. The side effect inside the step runs at least once. Use idempotency keys or overwrite semantics.
- Using the orchestrator as a data plane. XCom, workflow payloads and state machine inputs are for references and small values. Pass object-store paths, not DataFrames.
- Ignoring parse and replay costs. Airflow re-parses Dag files continuously; durable engines replay history on cache misses. Heavy top-level code or huge histories slow everything.
- Picking by popularity instead of by the shape of a run. Write the one-sentence description of a run first.
The cluster map
- DAG Fundamentals & Airflow Architecture - Topological Scheduling, the Scheduler Loop, Executors, Task States & Pools: why DAGs, topological waves and the critical path, Airflow 3.x components, task-instance states, pools and concurrency knobs, executors and HA schedulers.
- Writing Correct Pipelines - Data Intervals, Catchup & Backfill, Idempotent Tasks, Deferrable Sensors & Assets: data intervals, catchup vs backfill, idempotent partition overwrites, DST, XCom limits, sensors vs deferrable operators vs assets, and Deadline Alerts.
- Durable Execution & the Temporal Model - Workflows vs Activities, Event History, Replay & Determinism: workflows vs activities, event history and replay, determinism rules, timers, signals, queries and updates, the four activity timeouts, and what exactly-once means.
- Temporal Patterns in Production - Sagas, Child Workflows, Continue-As-New, Versioning & Worker Scaling: sagas with compensations, child workflows, continue-as-new and history limits, patching vs Worker Versioning, human waits, sticky queues and history shards.
- Choosing & Operating Workflow Orchestrators - Airflow vs Dagster vs Prefect vs Temporal vs Step Functions vs Argo, Testing & Migration: Airflow vs Dagster vs Prefect vs Argo vs Temporal vs Step Functions vs cron plus a queue, observability, DAG integrity and replay tests, and migrations.
Interview Q&A
What do Airflow and Temporal have in common?
Answer
Both persist state outside workers, dispatch work through queues to worker processes, retry with policies, and detect lost workers with heartbeats or timeouts. Both therefore deliver at-least-once execution of each unit of work, so units must be idempotent.
What is the core difference between a DAG scheduler and a durable execution engine?
Answer
A DAG scheduler runs a declared graph per schedule interval or data event and tracks state per task, retrying a failed task from its start. A durable engine runs ordinary code per business event and journals every step, so after a crash it replays the code and skips completed steps.
Why is a static graph an advantage for data pipelines?
Answer
The scheduler can show the whole plan before running, compute what is runnable, backfill any date range, rerun a single task for a single partition, and draw lineage. Static structure is what makes catchup and backfill cheap.
When would you choose Step Functions over Temporal?
Answer
When you are all-in on AWS, the flow is mostly calls to AWS services, you want zero servers to operate, and the logic fits a state machine. Choose Temporal when the logic is code-heavy, you need multi-cloud or self-hosting, or you want rich signals, queries, updates and versioning in code.
Is a durable workflow just a saga?
Answer
A saga is a pattern (a sequence of local transactions with compensations). A durable workflow engine is a good place to implement an orchestrated saga because the compensation logic itself is guaranteed to finish.
What does a heartbeat timeout detect that a run timeout does not?
Answer
A worker that died silently on a long task. Without heartbeats you only notice when the whole attempt times out.
Where do Restate and Inngest fit?
Answer
Both are durable execution engines: Restate journals handlers in ordinary services, Inngest checkpoints step.run blocks inside functions.
Why should big data not flow through the orchestrator?
Answer
Its state store is what scheduling depends on. Pass object-store paths and keep payloads small.
Check yourself
Pick two multi-step jobs you run today. Write one sentence for what a run is in each, decide which family fits, and list the side effect in each that needs an idempotency key.
Elsewhere in the library
These pages stay as they are. This lesson only points at them: Sagas & Distributed Transactions — Orchestration, Choreography & Compensations, Orchestration vs Choreography — Central Coordinator vs Event Dance, Dependency Graphs & Incremental Builds - Task DAGs, Topological Scheduling, Content Hashing & Affected Detection, Cloudflare Durable Objects - Single-Instance Actors, Edge State & When to Use Them, Stream Processing — Event Time, Windows, State & Exactly-Once, Event-Driven Architecture — Sync vs Events, Patterns & Tradeoffs, At-Least-Once vs Exactly-Once Delivery.