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.