Lesson 14 / 25
Sensors and Deferrable Operators
Wait for files, partitions or external events without wasting worker slots.
Waiting efficiently
A sensor is a task that waits until a condition is true: a file appears in object storage, a table partition exists, an external API returns "ready". In poke mode it holds a worker slot while it checks repeatedly; in reschedule mode it frees the slot between checks. Best of all, many sensors are deferrable: they hand the waiting to a lightweight triggerer process using asyncio, using almost no worker resources, so thousands can wait at once. Always set a timeout so a sensor that never succeeds eventually fails instead of waiting forever.
A sensor with a timeout
Illustrative; needs the Amazon provider and an AWS connection. Waits up to 6 hours, checking every 5 minutes without holding a worker.
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
wait_for_file = S3KeySensor(
task_id="wait_for_file",
bucket_name="lake",
bucket_key="incoming/orders_{{ ds }}.csv",
aws_conn_id="aws_default",
poke_interval=300, # seconds between checks
timeout=6 * 60 * 60, # give up after 6 hours
mode="reschedule", # free the worker slot between checks
# deferrable=True, # even lighter: uses the triggerer
)Prefer event-driven when possible
If the producer can notify you (for example by triggering your DAG or updating an Asset), that is more efficient and faster than polling. Use sensors when you can only observe the world.
Quick check: Why set a timeout on a sensor?
- It deletes the DAG
- Timeouts make sensors faster
- Sensors require it syntactically
- So it fails instead of waiting forever if the condition never becomes true
Answer
So it fails instead of waiting forever if the condition never becomes true — Bounded waiting prevents silent, indefinite hangs.