To react to the outside world rather than Airflow 129 's own tasks, Airflow 3 adds asset watchers: an AssetWatcher attaches a trigger (a BaseEventTrigger) to an asset, the triggerer runs it continuously, and each event becomes an asset event. The common-messaging provider's MessageQueueTrigger covers queues such as Amazon SQS, Kafka 129 and Redis 2,763 pub/sub; for BookNest's export, FileDeleteTrigger waits for a marker file and deletes it:
from airflow.providers.standard.triggers.file import FileDeleteTrigger
from airflow.sdk import Asset, AssetWatcher, dag, task
READY = "/opt/airflow/data/landing/_READY"
export_ready = Asset(name="order_export_ready", watchers=[
AssetWatcher(name="ready_file", trigger=FileDeleteTrigger(filepath=READY))])
@dag(schedule=[export_ready], tags=["booknest"])
def booknest_on_drop():
@task
def react(triggering_asset_events=None):
for asset, events in triggering_asset_events.items():
print(f"{asset.name}: {len(events)} event(s), extra={events[0].extra}")
react()
booknest_on_drop()Once the DAG was unpaused, the triggerer ran the watcher; touching the file started a run:
20:58:06 touch landing/_READY
20:58:07 file gone
booknest_on_drop asset_triggered__2026-10-01T20:58:07.713942+00:00_Dc7rhFyB success
order_export_ready: 1 event(s), extra={'from_trigger': True, 'payload': True}No worker slot was held while waiting. Watchers suit low-volume signals ("a batch is ready"), not one run per record: consume order streams with Kafka (Apache Kafka and Managed Cloud Kafka).