Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Assets and data-aware scheduling

Airflow & DAGs · Airflow 3

Assets and data-aware scheduling

Mediumairflow-42
airflow-3assetsdata-aware-schedulingevent-driven

Question

What are Assets (formerly Datasets), and how do they trigger DAGs?

Solution

Assets (introduced as Datasets in Airflow 2.4 and renamed in Airflow 3) allow DAGs to trigger automatically when underlying data updates, shifting workflows from rigid clock schedules to event-driven pipelines.

How producer and consumer tasks connect

An Asset is defined by a URI string that represents a logical data entity. Airflow does not connect to the URI or poll external files. Instead, it tracks metadata events emitted by producer tasks.

from airflow.sdk import Asset, DAG, task

orders_asset = Asset("s3://analytics-bucket/raw/orders")

with DAG(dag_id="orders_ingest", schedule="@hourly"):
    @task(outlets=[orders_asset])
    def extract_orders():
        # writes files to s3
        return "done"

with DAG(dag_id="orders_transform", schedule=[orders_asset]):
    @task
    def aggregate_orders():
        # runs as soon as extract_orders succeeds
        return "aggregated"

The consumer DAG triggers only after the producer task completes successfully. If extract_orders fails, no asset update event is recorded, and orders_transform never starts.

Combining assets with Boolean logic

Consumer DAGs can combine multiple assets using logical operators:

schedule = (Asset("s3://bucket/orders") & Asset("s3://bucket/users")) | Asset("s3://bucket/overrides")

The pipeline triggers when both orders and users are updated, or whenever an emergency overrides asset is written.

Airflow 3 asset watchers

Airflow 3 adds the @asset decorator, allowing developers to define an asset and its producing logic together in a single Python function. It also introduces asset watchers, which monitor external event infrastructure like Kafka topics, AWS SQS queues, or Google Cloud Pub/Sub subscriptions. When an external message lands in the queue, the asset watcher records an asset update event in Airflow, triggering downstream DAGs without needing continuous polling sensors.

PreviousNext