In Airflow, when does a DAG scheduled on Datasets actually run?
answer
- the clock is replaced by a declared dependency
- the producer's success is the signal
- all listed datasets, since the last run
- the URI is a name, not a location Airflow checks
- a stalled producer means silence, not an error
basics
~20 sIt runs when every Dataset listed in its schedule has been updated at least once since its last run. A Dataset is marked updated when a task that declares it in outlets finishes successfully — Airflow never inspects the underlying storage.
solid answer
~50 sData-driven scheduling replaces the clock with a declared dependency. A producing task lists `outlets=[Dataset("s3://warehouse/orders")]`; when that task instance **succeeds**, Airflow records a dataset event. A consuming DAG declares `schedule=[orders, customers]`, and the scheduler creates a run for it once *all* listed datasets have received at least one event since the consumer's previous run — plain AND semantics. From Airflow 2.9 you can express OR and mixed logic with `DatasetAny`/`DatasetAll`, and combine a dataset condition with a time schedule. Two things routinely surprise people: the dataset URI is only an **identifier** — Airflow never checks S3 or the warehouse, so a task that succeeds while writing nothing still fires the event — and a dataset-triggered run has no meaningful data interval, so partition logic must not derive a window from it. Airflow 3 renames Datasets to Assets.
code
python · 15 linesimport pendulum
from airflow import DAG
from airflow.datasets import Dataset
from airflow.operators.python import PythonOperator
orders = Dataset("s3://warehouse/orders")
with DAG("ingest_orders", schedule="@daily",
start_date=pendulum.datetime(2024, 5, 1, tz="UTC"), catchup=False):
PythonOperator(task_id="write_orders", python_callable=write,
outlets=[orders]) # event emitted on success
with DAG("orders_mart", schedule=[orders], # no clock at all
start_date=pendulum.datetime(2024, 5, 1, tz="UTC"), catchup=False):
PythonOperator(task_id="build_mart", python_callable=build)go deeper
Know that a DAG can be scheduled on data rather than a clock, and that the trigger is a producing task's successful completion rather than a file appearing in storage.
Explain the mechanics precisely: outlets on the producer, a list in the consumer's schedule, conjunctive semantics, and the fact that the URI is only an identifier Airflow never inspects.
Show the operational judgment — alerting on absence rather than failure, choosing which task carries the outlet so its success means something, and knowing that dataset-triggered runs have no interval to base partition logic on.
Own the choice between time-driven and data-driven scheduling across teams: where declared data contracts beat cron offsets, where they hide stalls, and how replay and cross-deployment dependencies are handled when the mechanism cannot express them.
## The idea Time-based scheduling encodes a *guess*: "upstream is usually done by 03:00, so start at 03:30". Data-driven scheduling encodes the *actual* dependency: "start when the orders table has been refreshed". Airflow 2.4 introduced Datasets for this; Airflow 3 renames them to Assets and extends them with external watchers. ## The mechanics A `Dataset` is constructed from a URI string: `Dataset("s3://warehouse/orders")`. The producer declares it as an outlet on a task; the consumer lists it in `schedule`. ```python from airflow.datasets import Dataset orders = Dataset("s3://warehouse/orders") # producer with DAG("ingest_orders", schedule="@daily", start_date=..., catchup=False): PythonOperator(task_id="write_orders", python_callable=write, outlets=[orders]) # consumer — no schedule of its own with DAG("orders_mart", schedule=[orders], start_date=..., catchup=False): PythonOperator(task_id="build_mart", python_callable=build) ``` When `write_orders` reaches the `success` state, Airflow writes a **dataset event** against that URI. The scheduler watches consumers: `orders_mart` becomes eligible once `orders` has at least one event newer than the consumer's last run. With several datasets in the list, the default is conjunctive — *all* of them must have been updated since the last run, so a consumer of a daily and a weekly feed runs weekly, not daily. Airflow 2.9 added `DatasetAny(...)` and `DatasetAll(...)` for explicit OR/AND expressions and a combined time-or-dataset schedule for "run when data arrives, but at least once a day regardless". ## What the URI is and is not The URI is a **name**, not a location Airflow polls. Nothing verifies that anything was written to `s3://warehouse/orders`. The consequences follow directly: - A producing task that succeeds without writing a byte still emits the event, and the consumer runs on stale data. - A task that writes the data but then fails on a later line emits **no** event, and the consumer silently never runs — a stall with no error anywhere, which is the failure mode teams find hardest to notice. - Producer and consumer must spell the URI identically. A trailing slash or a differing scheme creates two unrelated datasets and a dependency that quietly does not exist. Because the semantics are "a task succeeded", the granularity of the guarantee is the granularity of the task. Make the outlet-bearing task the one whose success genuinely means "the data is complete and readable". ## Scheduling implications **No data interval.** A dataset-triggered run is not derived from a timetable, so there is no window it represents. Do not build `WHERE created_at >= data_interval_start` logic in a dataset-scheduled DAG; pass the window explicitly (through the triggering run's own knowledge of it, or by having the consumer read a watermark) or design the consumer as a full refresh of whatever is currently present. **No catchup.** There is nothing to catch up: dataset events are point-in-time facts, not a sequence of intervals. If you need to replay history through a dataset-scheduled DAG, you trigger it manually or backfill it as a DAG with an explicit range — the dataset mechanism will not reconstruct a past sequence for you. **Event coalescing.** If a dataset is updated several times while the consumer is already running, the consumer does not run once per event; it becomes eligible again after it finishes. Consumers should therefore be written to process "everything new", not "exactly one upstream batch". **Deployment scope.** Datasets link DAGs inside one Airflow deployment. Two separate Airflow instances do not share dataset events by that mechanism alone. ## When to reach for it Datasets are the right tool when the real dependency is data readiness and the producers' timing is variable — an upstream job that sometimes finishes at 03:10 and sometimes at 06:40. They remove the brittle offset guessing that a cron-plus-sensor arrangement encodes, they show the dependency in the UI's dataset graph, and they let a producer's schedule change without every consumer needing an edit. They are the wrong tool when the consumer needs a defined data window per run (an interval-partitioned incremental load), when the dependency crosses Airflow deployments, or when you need historical replay through the same path, since dataset triggering has no notion of an interval to replay. ## Operational advice Monitor for **absence**: a dataset-scheduled DAG that has not run in longer than its expected cadence is the signal that an upstream task failed before emitting its event. Because nothing is scheduled, there is no failed run to alert on — the DAG simply goes quiet. That is the single most important alert to add when a team adopts data-driven scheduling.
- A dataset-scheduled Airflow DAG has not run for three days and no run has failed. What happened?Almost certainly an upstream producing task never reached success, so no dataset event was recorded and the consumer was never made eligible. Because nothing was scheduled, there is no failed consumer run to alert on. Check the producer's recent task instances, verify the outlet URI matches the consumer's exactly, and add a freshness alert on 'DAG has not run in longer than expected'.
- If a DAG lists two Datasets in its schedule and one is updated ten times while the other is not updated at all, does it run?No. The default semantics are conjunctive: every dataset in the list must have at least one event since the consumer's last run. Ten events on one dataset do not compensate for zero on the other, and they do not queue ten runs either. Airflow 2.9 and later let you express OR explicitly with DatasetAny if that is what you want.
- Why should a dataset-triggered Airflow DAG not filter source data on data_interval_start?Because the run was not derived from a timetable, so it does not represent a data window. Any interval values it carries reflect the trigger moment rather than a meaningful range, and partition logic built on them silently reads the wrong rows. Pass the window explicitly or design the consumer to process everything new since its own recorded watermark.
- Does Airflow verify that data was actually written to the dataset URI?No. The URI is purely an identifier used to match producers to consumers. The event is recorded because a task declaring that outlet reached the success state, whatever it did or did not write. If the guarantee matters, make the outlet-bearing task one whose success genuinely means the data is complete, and add a data-quality check ahead of it.
saying these in an interview costs you the question
- Thinks Airflow polls the dataset URI to detect a file appearing
- Assumes a consumer of several datasets runs when any one of them updates
- Believes several updates during a consumer run queue one consumer run each
- Builds interval-partition filters inside a dataset-triggered DAG
- Expects catchup to replay historical dataset events