Event-Driven Scheduling

Event-Driven Scheduling and Asset Watchers

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:

booknest_on_drop.py: run whenever the order service drops a _READY filePython
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:

Output of 9
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).