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.