Data engineering
Part 3 of 6 · Workflow OrchestrationWriting Correct Pipelines - Data Intervals, Catchup & Backfill, Idempotent Tasks, Deferrable Sensors & Assets
Logical date and data intervals, catchup and backfill in Airflow 3 (catchup off by default, logical_date None for asset/API runs), idempotent tasks with partition overwrite (runnable sqlite demo: append vs delete+insert), DST and time zones (runnable America/Chicago demo), XCom limits, sensors vs deferrable operators vs assets and AssetWatcher, dynamic task mapping, Deadline Alerts replacing SLAs, failure modes; Dagster/Prefect comparisons; waiting decision chart.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
Why does the daily run for October 5 start on October 6?
Answer
It processes [Oct 5, Oct 6), which is only complete after the interval ends.
L2
What makes a task idempotent?
Answer
Running it once or many times for the same interval leaves the same final state.
L3
Catchup vs backfill?
Answer
Catchup is the scheduler creating missed runs; backfill is a chosen date range started by a person or automation.
L4
What changed for intervals in Airflow 3?
Answer
Asset and API runs can have logical_date None and no data interval, and catchup is off by default.
L5
Why avoid scheduling at 02:30 local in a DST zone?
Answer
That time is skipped once a year and repeated once a year.
L6
Sensor or deferrable operator?
Answer
Deferrable for any meaningful wait; it frees the worker slot.
L7
What replaced SLAs in Airflow 3?
Answer
Deadline Alerts with a reference time, an interval and a callback.
Failure modes
Double-counted totals after a retry
An appending task inserted the same interval's rows twice.
Backfill recomputes the wrong day
The task derived its window from the wall clock, so every backfill run processed yesterday.
Partition count checks fail on DST days
Local days of 23 or 25 hours break a hard-coded 24 hourly partitions.
Misconceptions
The logical date is when the run executes.
It is the start of the data interval.
Every run has a data interval.
Asset-triggered and API-triggered runs in Airflow 3 may not.
XCom is a fine way to pass a DataFrame.
XCom lives in the metadata DB; pass a path to object storage.
Interviewer traps
Using WHERE ts > now() - 1 day in a task.
Read data_interval_start and data_interval_end from the run context.
Turning catchup on for a long-paused Dag without limits.
Cap max_active_runs and use pools so the source is not flooded.
Design scenario
Same prompt for every reader.
Requirements
Idempotent daily partitions, safe backfills of the last 90 days, late-data corrections, and an alert when the partition is late.
Failure assumptions
- A worker loses its ack and the task retries.
- Events arrive up to two days late.
- The DST change shortens one day.
Constraints
- Source is a shared OLTP database.
- Warehouse supports partition overwrite.
Prompt
Design a daily revenue pipeline that must be rerunnable for any day, tolerate late events and publish a dashboard partition by 06:30 Chicago time.
API
Which context fields does each task read, and how are reruns triggered?
Data
How are partitions keyed and overwritten, and how are late events folded in?
Architecture
Where do sensors, deferrable operators, assets and Deadline Alerts fit?
Overview
A scheduled pipeline is correct when running any interval again, for any reason, produces the same result. Orchestrators will retry tasks, catch up missed intervals, backfill old ones, and occasionally run something twice because a worker looked dead. So each task should be a pure function of its data interval: read only the slice [data_interval_start, data_interval_end), write only the partition for that slice, and overwrite rather than append. Everything else on this page (logical dates, catchup, XCom limits, deferrable sensors, dynamic mapping, assets, deadline alerts, DST) is a consequence of that one rule or a guard against breaking it. Dagster's asset model and Prefect's caching reach the same goal from different directions.
Logical date and the data interval
For a time-scheduled run, Airflow passes the run a data interval: a daily run "for 2026-10-05" covers [2026-10-05 00:00, 2026-10-06 00:00) and actually starts after the interval ends, because only then is the data complete. The logical date (formerly execution_date, removed in 3.0) is the interval start. This confuses everyone once: the run labelled 2026-10-05 executes on 2026-10-06.
Two 3.x changes matter:
| Change (Airflow 3.0) | Consequence |
|---|---|
Runs triggered by an asset event, or through the REST API without a logical date, get logical_date=None and no data interval | Tasks that read data_interval_start from context get a KeyError; write tasks that handle both trigger types |
create_cron_data_intervals now defaults to False, so a bare cron string uses CronTriggerTimetable | A cron-scheduled run is a point-in-time trigger unless you opt back into interval semantics; decide explicitly before relying on ds or data_interval_* |
catchup_by_default is False | New Dags no longer create a run for every missed interval on activation; set catchup=True deliberately |
| Future logical dates are rejected | Use logical_date=None for "run now" triggers |
Append rows or overwrite the interval's partition?
Prefer
Delete and insert the interval's partition
Read only the run's interval and replace its partition in one transaction.
- A retry left one row for 2026-10-05, not two.
- The backfill produced correct totals for 2026-10-03 and 2026-10-04.
- A late event was corrected in place with no duplicates.
Alternative
Append rows from a wall-clock window
Insert whatever 'yesterday' means when the task happens to run.
- A retry duplicated the row and doubled the totals.
- Every backfill run recomputed 2026-10-05.
- Nobody notices until a reconciliation.
A correct daily partition run
Diagram 1 condensed: interval in, partition replaced, failure path out.
- 1
Receive the interval
The run gets [start, end) for its logical date. - 2
Wait without a slot
A deferrable operator or asset trigger waits for the input. - 3
Read only the slice
The task filters input by the interval, not by now(). - 4
Overwrite the partition
Delete and insert, or a partition swap, in one transaction. - 5
Retry or backfill safely
Reruns replace the same partition instead of adding rows.
Catchup vs backfill vs rerun
| Operation | Who starts it | What it does | Risk |
|---|---|---|---|
| Catchup | The scheduler, when catchup=True and intervals were missed | Creates a run for every missed interval since the last one | A paused Dag unpaused after a month floods workers; cap with max_active_runs |
| Backfill | A human or automation (UI, REST API, airflow backfill create) | Creates runs for a chosen date range; in 3.x backfills are scheduler-managed and visible like normal runs | Overwrites history: only safe if tasks are idempotent |
| Clear / rerun | A human, per task or per run | Resets task instances so they run again | Same as backfill; with Dag versioning, choose whether reruns use the original or latest code (rerun_with_latest_version, 3.3) |
Idempotent and deterministic tasks
"""Why pipeline tasks must be idempotent and interval-driven.
A daily revenue job is run three ways against SQLite:
BAD : appends with INSERT and uses the wall clock ("yesterday") to pick rows
GOOD : reads ONLY its data interval [start, end) and overwrites its partition
(DELETE + INSERT for that day inside one transaction)
Then we do what real orchestrators do: retry a task, and backfill old dates.
"""
import sqlite3
db = sqlite3.connect(":memory:")
db.executescript("""
CREATE TABLE events(ts TEXT, amount INTEGER);
CREATE TABLE revenue_bad(day TEXT, total INTEGER);
CREATE TABLE revenue_good(day TEXT PRIMARY KEY, total INTEGER);
""")
db.executemany("INSERT INTO events VALUES (?,?)", [
("2026-10-03T09:00", 100), ("2026-10-03T18:30", 50),
("2026-10-04T11:00", 70),
("2026-10-05T08:15", 30), ("2026-10-05T23:59", 20),
])
WALL_CLOCK_TODAY = "2026-10-06" # the day the job happens to execute
def bad_task():
# Picks "yesterday" from the wall clock and appends. Ignores the run's interval.
db.execute("""INSERT INTO revenue_bad
SELECT date(?, '-1 day'), COALESCE(SUM(amount),0) FROM events
WHERE ts >= date(?, '-1 day') AND ts < ?""",
(WALL_CLOCK_TODAY, WALL_CLOCK_TODAY, WALL_CLOCK_TODAY))
def good_task(start, end):
# Deterministic: same interval in -> same partition out, however often it runs.
with db: # one transaction: readers never see a half-written partition
db.execute("DELETE FROM revenue_good WHERE day = ?", (start[:10],))
db.execute("""INSERT INTO revenue_good
SELECT ?, COALESCE(SUM(amount),0) FROM events
WHERE ts >= ? AND ts < ?""", (start[:10], start, end))
def show(table):
return db.execute(f"SELECT day, total FROM {table} ORDER BY day").fetchall()
print("1) Normal run for interval 2026-10-05, then the task is retried (worker lost its ack)")
bad_task(); bad_task()
good_task("2026-10-05T00:00", "2026-10-06T00:00"); good_task("2026-10-05T00:00", "2026-10-06T00:00")
print(" BAD :", show("revenue_bad"), "<- duplicated row, totals now double-count")
print(" GOOD:", show("revenue_good"))
print("2) Backfill 2026-10-03 .. 2026-10-04 (the scheduler passes each run its own interval)")
for day, nxt in [("2026-10-03", "2026-10-04"), ("2026-10-04", "2026-10-05")]:
bad_task()
good_task(day + "T00:00", nxt + "T00:00")
print(" BAD :", show("revenue_bad"), "<- every backfill run recomputed 'yesterday'")
print(" GOOD:", show("revenue_good"))
print("3) A late event for 2026-10-04 arrives; rerun just that interval")
db.execute("INSERT INTO events VALUES ('2026-10-04T13:00', 5)")
good_task("2026-10-04T00:00", "2026-10-05T00:00")
print(" GOOD:", show("revenue_good"), "<- corrected in place, no duplicates")Output:
1) Normal run for interval 2026-10-05, then the task is retried (worker lost its ack)
BAD : [('2026-10-05', 50), ('2026-10-05', 50)] <- duplicated row, totals now double-count
GOOD: [('2026-10-05', 50)]
2) Backfill 2026-10-03 .. 2026-10-04 (the scheduler passes each run its own interval)
BAD : [('2026-10-05', 50), ('2026-10-05', 50), ('2026-10-05', 50), ('2026-10-05', 50)] <- every backfill run recomputed 'yesterday'
GOOD: [('2026-10-03', 150), ('2026-10-04', 70), ('2026-10-05', 50)]
3) A late event for 2026-10-04 arrives; rerun just that interval
GOOD: [('2026-10-03', 150), ('2026-10-04', 75), ('2026-10-05', 50)] <- corrected in place, no duplicatesThe bad task commits two sins at once: it appends (not idempotent under retry) and it derives its window from the wall clock (not deterministic under backfill). The good task is what Maxime Beauchemin (Airflow's creator) called functional data engineering: tasks are pure functions from an immutable input partition to an output partition, and outputs are replaced, never mutated.
| Write pattern | Idempotent? | Notes |
|---|---|---|
INSERT INTO t SELECT ... | No | Every retry adds rows |
DELETE WHERE partition = x then INSERT in one transaction | Yes | Works on any SQL database |
INSERT OVERWRITE ... PARTITION (dt = x) / partition swap | Yes | Native in Hive, Spark, BigQuery partition decorators, Iceberg and Delta |
MERGE / upsert on a natural key | Yes, for unchanged keys | Late deletes need extra handling |
| Write to a temp path, then atomic rename or pointer swap | Yes | Object stores: write a new prefix, then flip a manifest or table pointer |
| Calling an external API (send email, charge card) | Only with an idempotency key | Or move it out of the batch DAG into a durable workflow |
Timezones and DST
Schedules are written in local time by humans and executed in UTC by machines. The gap between the two produces real outages:
// Timezone and DST bugs in scheduled pipelines, computed with the real tz database (Intl).
// Zone: America/Chicago. In 2026 US DST starts Sun 2026-03-08 02:00 local (clocks jump to 03:00)
// and ends Sun 2026-11-01 02:00 local (clocks fall back to 01:00, so 01:00-01:59 happens twice).
const ZONE = "America/Chicago";
const fmt = new Intl.DateTimeFormat("en-US", {
timeZone: ZONE, hourCycle: "h23", year: "numeric", month: "2-digit", day: "2-digit", hour: "2-digit", minute: "2-digit",
});
function wall(ms: number): string {
const p = Object.fromEntries(fmt.formatToParts(new Date(ms)).map((x) => [x.type, x.value]));
return `${p.year}-${p.month}-${p.day} ${p.hour}:${p.minute}`;
}
// All UTC instants whose local wall clock equals the requested time (0 = gap, 2 = overlap).
function localToUtc(date: string, hhmm: string): number[] {
const [y, m, d] = date.split("-").map(Number); const [hh, mm] = hhmm.split(":").map(Number);
const naive = Date.UTC(y, m - 1, d, hh, mm);
const hits = [-6, -5].map((off) => naive - off * 3600_000).filter((ms) => wall(ms) === `${date} ${hhmm}`);
return [...new Set(hits)].sort();
}
const iso = (ms: number) => new Date(ms).toISOString().slice(0, 16) + "Z";
console.log("1) A daily job at a local wall-clock time");
for (const [date, t] of [["2026-03-07", "02:30"], ["2026-03-08", "02:30"], ["2026-11-01", "01:30"], ["2026-11-02", "01:30"]]) {
const hits = localToUtc(date, t);
const verdict = hits.length === 0 ? "DOES NOT EXIST (spring-forward gap): a naive scheduler may skip it"
: hits.length === 2 ? "HAPPENS TWICE (fall-back overlap): a naive scheduler may run it twice" : "ok";
console.log(` ${date} ${t} local -> ${hits.map(iso).join(", ") || "-"} ${verdict}`);
}
console.log("2) A job pinned to 08:30 UTC, written as if it were always 02:30 local");
for (const date of ["2026-03-07", "2026-03-09", "2026-11-02"]) {
const [y, m, d] = date.split("-").map(Number);
console.log(` ${date} 08:30Z runs at local ${wall(Date.UTC(y, m - 1, d, 8, 30)).slice(11)}`);
}
console.log("3) Local-day data intervals are not always 24 hours");
for (const [day, next] of [["2026-03-07", "2026-03-08"], ["2026-03-08", "2026-03-09"], ["2026-11-01", "2026-11-02"]]) {
const hours = (localToUtc(next, "00:00")[0] - localToUtc(day, "00:00")[0]) / 3600_000;
const check = hours === 24 ? "ok" : `a check expecting 24 hourly partitions FAILS (found ${hours})`;
console.log(` [${day} 00:00, ${next} 00:00) local = ${hours}h ${check}`);
}
console.log("Fix: store and compare instants in UTC, define intervals with a tz-aware timetable,");
console.log(" never schedule inside 01:00-03:00 local, and derive partition counts from the interval.");Output:
1) A daily job at a local wall-clock time
2026-03-07 02:30 local -> 2026-03-07T08:30Z ok
2026-03-08 02:30 local -> - DOES NOT EXIST (spring-forward gap): a naive scheduler may skip it
2026-11-01 01:30 local -> 2026-11-01T06:30Z, 2026-11-01T07:30Z HAPPENS TWICE (fall-back overlap): a naive scheduler may run it twice
2026-11-02 01:30 local -> 2026-11-02T07:30Z ok
2) A job pinned to 08:30 UTC, written as if it were always 02:30 local
2026-03-07 08:30Z runs at local 02:30
2026-03-09 08:30Z runs at local 03:30
2026-11-02 08:30Z runs at local 02:30
3) Local-day data intervals are not always 24 hours
[2026-03-07 00:00, 2026-03-08 00:00) local = 24h ok
[2026-03-08 00:00, 2026-03-09 00:00) local = 23h a check expecting 24 hourly partitions FAILS (found 23)
[2026-11-01 00:00, 2026-11-02 00:00) local = 25h a check expecting 24 hourly partitions FAILS (found 25)
Fix: store and compare instants in UTC, define intervals with a tz-aware timetable,
never schedule inside 01:00-03:00 local, and derive partition counts from the interval.Expected1) A daily job at a local wall-clock time 2026-03-07 02:30 local -> 2026-03-07T08:30Z ok 2026-03-08 02:30 local -> - DOES NOT EXIST (spring-forward gap): a naive scheduler may skip it 2026-11-01 01:30 local -> 2026-11-01T06:30Z, 2026-11-01T07:30Z HAPPENS TWICE (fall-back overlap): a naive scheduler may run it twice 2026-11-02 01:30 local -> 2026-11-02T07:30Z ok 2) A job pinned to 08:30 UTC, written as if it were always 02:30 local 2026-03-07 08:30Z runs at local 02:30 2026-03-09 08:30Z runs at local 03:30 2026-11-02 08:30Z runs at local 02:30 3) Local-day data intervals are not always 24 hours [2026-03-07 00:00, 2026-03-08 00:00) local = 24h ok [2026-03-08 00:00, 2026-03-09 00:00) local = 23h a check expecting 24 hourly partitions FAILS (found 23) [2026-11-01 00:00, 2026-11-02 00:00) local = 25h a check expecting 24 hourly partitions FAILS (found 25) Fix: store and compare instants in UTC, define intervals with a tz-aware timetable, never schedule inside 01:00-03:00 local, and derive partition counts from the interval.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Airflow stores datetimes in UTC and supports timezone-aware schedules via the Dag's start_date timezone; still, the safe habits are: keep start_date timezone-aware, avoid scheduling between 01:00 and 03:00 local in DST zones, and never hard-code "24 hourly partitions per day".
XCom: metadata, not data
XComs move small values between tasks and are stored in the metadata database by default. The docs say it plainly: designed for small amounts of data, do not pass large values like DataFrames. Large payloads bloat the database that the scheduler itself depends on. Pass a path to object storage instead, or configure an object-storage XCom backend so large values are offloaded automatically. Every orchestrator has the same rule: Step Functions caps state input and output at 256 KiB, and Temporal warns at 256 KB per payload and errors at 2 MB.
Waiting for things: sensors, reschedule, deferrable, assets
| Mechanism | How it waits | Cost while waiting | Use when |
|---|---|---|---|
Sensor, poke mode | Holds a worker slot and polls in a loop | One worker slot per waiting sensor | Very short waits only |
Sensor, reschedule mode | Frees the slot between pokes (up_for_reschedule) | Scheduler overhead per poke | Waits of minutes to hours, legacy sensors |
| Deferrable operator or trigger | Defers to an async trigger in the triggerer | No worker slot; by default no pool slot | Long waits on external systems; thousands of waits |
| Asset (data-aware) scheduling | The consumer Dag is triggered when producers update an asset | Nothing: no waiting task exists | You control the producer, or an asset watcher can receive the event |
| Event-driven asset watchers | A trigger listens to a message queue and emits asset events | One trigger in the triggerer | Start Dags from external events (for example, a queue message) |
Assets flip the question from "poll until the file appears" to "run when the upstream asset changes". Airflow 3.0 renamed datasets to assets and added the @asset decorator; 3.2 added asset partitions, so a downstream run is triggered by the specific partition that changed, and 3.3 added partition mappers for rollups and fan-out.
Diagram 1: a correct daily partition run, including the failure path
Flow
- 1
Step 1: scheduler creates the run for interval 2026-10-05 00:00 to 2026-10-06 00:00 UTC
- nextStep 2: wait for input via an asset event or a deferrable trigger, no worker slot held
- 2
Step 2: wait for input via an asset event or a deferrable trigger, no worker slot held
- nextStep 3: extract reads only rows inside the interval
- 3
Step 3: extract reads only rows inside the interval
- nextStep 4: transform writes to a staging location keyed by the interval
- 4
Step 4: transform writes to a staging location keyed by the interval
- nextStep 5: load replaces the 2026-10-05 partition in one transaction or atomic swap
- 5
Step 5: load replaces the 2026-10-05 partition in one transaction or atomic swap
- nextStep 6: data quality checks on the new partition
- 6
Step 6: data quality checks on the new partition
- nextStep 7: emit asset event so downstream Dags run for this partition
- nextFailure path: mark failed, alert, keep the previous partition visible
- 7
Step 7: emit asset event so downstream Dags run for this partition
- 8
Failure path: mark failed, alert, keep the previous partition visible
- nextFix the data or code, then clear the run: the same interval is recomputed and overwritten
- 9
Fix the data or code, then clear the run: the same interval is recomputed and overwritten
- nextStep 3: extract reads only rows inside the interval
Lesson map
Writing Correct Pipelines - Data Intervals, Catchup & Backfill, Idempotent Tasks, Deferrable Sensors & Assets
Logical date and data intervals, catchup and backfill in Airflow 3 (catchup off by default, logical_date None for asset/API runs), idempotent tasks with partition overwrite (runnable sqlite demo: append vs delete+insert), DST and time zones (runnable America/Chicago demo), XCom limits, sensors vs deferrable operators vs assets and AssetWatcher, dynamic task mapping, Deadline Alerts replacing SLAs, failure modes; Dagster/Prefect comparisons; waiting 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: scheduler creates the run for interval 2026-10-05 00:00 to 2026-10-06 00:00 UTC"] b["Step 2: wait for input via an asset event or a deferrable trigger, no worker slot held"] c["Step 3: extract reads only rows inside the interval"] d["Step 4: transform writes to a staging location keyed by the interval"] e["Step 5: load replaces the 2026-10-05 partition in one transaction or atomic swap"] f["Step 6: data quality checks on the new partition"] g["Step 7: emit asset event so downstream Dags run for this partition"] x["Failure path: mark failed, alert, keep the previous partition visible"] y["Fix the data or code, then clear the run: the same interval is recomputed and overwritten"] a -->|continues| b b -->|continues| c c -->|continues| d d -->|continues| e e -->|continues| f f -->|continues| g f -->|continues| x x -->|continues| y y -->|continues| c
Dynamic task mapping
Mapping turns one declared task into N task instances at run time (.expand() for a cross product, .expand_kwargs() for zipped arguments, .partial() for constants). Each mapped instance has its own state, retries and logs, so one bad file does not fail its siblings. Limits to remember: max_map_length (1024 by default) and the fact that a reducer downstream receives a lazy sequence of the mapped results, which still lives in XCom, so map over paths and keep results small.
Deadlines, SLAs and alerting
Airflow 2's sla and sla_miss_callback were removed in 3.0. Their replacement is Deadline Alerts (added in 3.1 as experimental, with synchronous callbacks in 3.2): you choose a reference point (queued time, logical date or a fixed datetime), an interval, and a callback. Alert on what consumers care about ("the 06:00 dashboard partition is not ready by 06:30"), not on every task failure, and pair it with per-task on_failure_callback notifiers for fast feedback.
Compared: Dagster assets and Prefect
| Concept | Airflow 3.x | Dagster | Prefect 3 |
|---|---|---|---|
| Primary abstraction | Dag of tasks; assets as a scheduling layer | Software-defined asset: code that describes data that should exist | Flow (a Python function) calling tasks |
| Partitions | Data interval per run; asset partitions since 3.2 | First-class partition definitions per asset: time, static, multi-dimensional, dynamic (docs recommend at most 100,000 partitions per asset) | Parameters; no built-in partition grid |
| Backfill unit | Date range of Dag runs | Set of asset partitions | Re-run flows with parameters |
| Data-aware triggering | Asset events, watchers | Declarative automation and sensors on asset state | Events and automations |
| Graph discovery | Declared, parsed ahead of time | Declared asset graph | Discovered as the Python runs |
| Idempotency helpers | Interval context, overwrite patterns | Partition key in context, IO managers | Task caching by inputs, transactions |
| Biggest strength | Ecosystem of providers, maturity | Lineage and asset health as first-class | Low ceremony, dynamic Python |
Diagram 2: how should a Dag wait?
Decisions
- 1
Step 1: the Dag needs something that is not ready yet
- nextStep 2: is it produced by another Airflow Dag?
- ?
Step 2: is it produced by another Airflow Dag?
- nextSchedule on the asset, the run is created when the producer emits
- nextStep 3: external system can push an event to a queue?
- 3
Schedule on the asset, the run is created when the producer emits
- ?
Step 3: external system can push an event to a queue?
- nextAsset with an AssetWatcher, event-driven scheduling
- nextStep 4: the wait can be minutes or hours?
- 5
Asset with an AssetWatcher, event-driven scheduling
- ?
Step 4: the wait can be minutes or hours?
- nextDeferrable operator or sensor, the triggerer polls, no worker slot held
- nextSensor in reschedule mode, or a short poke
- 7
Deferrable operator or sensor, the triggerer polls, no worker slot held
- nextStep 5: did it arrive before the timeout?
- 8
Sensor in reschedule mode, or a short poke
- ?
Step 5: did it arrive before the timeout?
- nextContinue with the interval's data
- nextFailure path: fail or soft-fail, alert, and decide whether to rerun the interval
- 10
Continue with the interval's data
- 11
Failure path: fail or soft-fail, alert, and decide whether to rerun the interval
What happens if you choose otherwise
- Wall-clock windows (
WHERE ts > now() - 1 day): fine until the first retry, backfill or late run, then silently wrong. - Appending writes: every retry double counts, and nobody notices until a finance reconciliation.
- Poke-mode sensors at scale: hundreds of sensors each pin a worker slot; real work starves. Use deferrable operators or assets.
- Catchup on with a long-paused Dag and no
max_active_runs: hundreds of runs hit the source database at once. - Keeping SLA code while upgrading to 3.x: it is gone; migrate to Deadline Alerts and failure callbacks.
Pitfalls
- Top-level code in Dag files (API calls,
Variable.get, heavy imports): runs on every parse, slows all Dags, and can take the Dag processor down. Move it inside tasks. - Non-deterministic task inputs:
datetime.now(), "latest" table snapshots, random sampling without a seed. - Dynamic
start_datesuch asdatetime.now(): the schedule never stabilizes. - Assuming every run has a data interval: asset-triggered and API-triggered runs have
logical_date=Nonein 3.x. - Ignoring late data: decide a lateness horizon and rerun recent partitions, or use the streaming tools where event time matters.
Interview Q&A
Why does the daily run for October 5 start on October 6?
Answer
The run processes the data interval [Oct 5, Oct 6), and that data is only complete once the interval ends. The logical date is the interval start, not the wall-clock execution time.
What makes a task idempotent, and why does the orchestrator need it?
Answer
Running it once or many times for the same interval leaves the same final state. Orchestrators retry, backfill and occasionally duplicate work after lost heartbeats, so only idempotent tasks make those operations safe. Overwrite the interval's partition instead of appending.
When do you use a deferrable operator instead of a sensor?
Answer
For any wait longer than a few seconds, especially at scale. Deferral moves the wait to an async trigger in the triggerer and frees the worker slot; a poke-mode sensor holds a slot the whole time.
Why not pass a DataFrame through XCom?
Answer
XComs live in the metadata database by default and are meant for small values. Large payloads bloat the database that scheduling depends on. Write to object storage and pass the path, or use an object-storage XCom backend.
How do Dagster assets change backfills?
Answer
You backfill partitions of assets ("rebuild these 30 daily partitions of orders_clean and everything downstream"), and Dagster plans the runs from the asset graph, instead of you choosing Dag runs by date.
What replaced Airflow SLAs in 3.x?
Answer
SLAs were removed in 3.0. Deadline Alerts (from 3.1) fire a callback when a Dag run has not finished within an interval of a reference time such as the logical date or queued time.
What does rerun_with_latest_version decide in Airflow 3.3?
Answer
Whether a cleared or rerun task uses the Dag code version its run started with or the latest version.
How do Dagster partitions change the backfill unit?
Answer
You backfill a set of asset partitions and Dagster plans runs from the asset graph instead of choosing Dag runs by date.
Check yourself
Take one of your scheduled jobs, rewrite its window as data_interval_start to data_interval_end, and change its write to delete-and-insert or a partition overwrite. Then rerun the same interval twice and compare row counts.
Elsewhere in the library
These pages stay as they are. This lesson only points at them: Backfills & Reconciliation — Idempotent Jobs, Checksums, Shadow Reads, API Idempotency Keys, Exactly-Once CDC Pipelines — Idempotent Consumers, Keys & At-Least-Once Reality, Event Time & Watermarks — Allowed Lateness, Side Outputs & Idle Sources, Partition Strategies — Range, Hash, List & Composite, Object Storage — Buckets, Keys, Consistency & Scale, Alerting — Multi-Window Burn Rates vs Static Thresholds.
Go Deeper
- Airflow: Dag Runs, catchup and backfill
- Airflow: Best Practices (top-level code, idempotency, testing)
- Airflow: XComs
- Airflow: Deferrable Operators and Triggers
- Airflow: Assets and asset partitions
- Airflow: Time zones
- Dagster: Partitioning assets
- Prefect: Tasks, caching and transactions
- Maxime Beauchemin: Functional Data Engineering