Menu

Airflow course · Lesson 3 of 5

Airflow Sensors and Deferrable Operators

Wait for files, other DAGs and external jobs without wasting workers: sensor modes, timeouts, ExternalTaskSensor, FileSensor, triggers and what replaced Smart Sensors.

  • Intermediate
  • 19 min read
  • Updated Oct 2026
On this page
  1. Running the examples
  2. Sensor basics
  3. What it is
  4. How it works
  5. Pitfalls
  6. In interviews
  7. Poke vs reschedule mode
  8. What it is
  9. How it works
  10. Pitfalls
  11. In interviews
  12. ExternalTaskSensor
  13. What it is
  14. How it works
  15. Pitfalls
  16. In interviews
  17. FileSensor
  18. What it is
  19. How it works
  20. Pitfalls
  21. In interviews
  22. Deferrable operators and triggers
  23. What it is
  24. How it works
  25. Pitfalls
  26. In interviews
  27. Smart Sensors (removed)
  28. What it was
  29. Why it went
  30. In interviews
  31. Practice questions
  32. Key takeaways

Pipelines spend a surprising amount of time waiting: for a file to land, for another team’s DAG to finish, for a Spark or warehouse job to complete. Sensors and deferrable operators are how Airflow waits, and choosing the wrong way to wait is one of the most common causes of an Airflow installation running out of worker slots. This lesson covers every waiting mechanism, from classic poking to triggers.

Running the examples

The examples use Airflow 3.3 with a metadata database created by airflow db migrate. The helper below is the same one used in the earlier lessons: it registers an in-memory DAG and runs it once with dag.test(). A useful detail for this lesson is that dag.test() also runs triggers itself, so deferrable tasks work without starting a separate airflow triggerer process.

import asyncio, contextlib, io, json, os, tempfile, time
from datetime import datetime, timedelta
import pendulum
from airflow.dag_processing.bundles.manager import DagBundlesManager
from airflow.dag_processing.dagbag import DagBag, sync_bag_to_db
from airflow.models.xcom import XComModel
from airflow.utils.session import create_session

def run_dag(dag, **test_kwargs):
    """Register an in-memory DAG, run it once with dag.test(), return (run state, task states, return values)."""
    with contextlib.redirect_stdout(io.StringIO()):
        DagBundlesManager().sync_bundles_to_db()
        bag = DagBag(dag_folder=tempfile.mkdtemp())
        bag.dags[dag.dag_id] = dag
        sync_bag_to_db(bag, "dags-folder", None)
        dr = dag.test(**test_kwargs)
    states, values = {}, {}
    with create_session() as session:
        for ti in dr.get_task_instances(session=session):
            key = ti.task_id if ti.map_index < 0 else f"{ti.task_id}[{ti.map_index}]"
            states[key] = str(ti.state)
        for x in session.query(XComModel).filter_by(dag_id=dag.dag_id, run_id=dr.run_id, key="return_value"):
            key = x.task_id if x.map_index < 0 else f"{x.task_id}[{x.map_index}]"
            values[key] = XComModel.deserialize_value(x)
    return str(dr.state), dict(sorted(states.items())), dict(sorted(values.items()))

# A landing folder with one file in it, exposed through a filesystem ("fs") connection
LANDING = tempfile.mkdtemp()
os.environ["AIRFLOW_CONN_FS_LANDING"] = json.dumps({"conn_type": "fs", "extra": {"path": LANDING}})
with open(os.path.join(LANDING, "orders_2026-03-01.csv"), "w") as fh:
    fh.write("order_id,amount\n1,40.0\n")
MARCH_1 = pendulum.datetime(2026, 3, 1, tz="UTC")

Sensor basics

What it is

A sensor is an operator whose only job is to wait until a condition is true: a file exists, a partition is present, an API reports “done”, another DAG finished. When the condition is met the sensor succeeds and downstream tasks run. If it is never met within the timeout, the sensor fails (or is skipped).

How it works

Every sensor subclasses BaseSensorOperator and implements poke(context), which returns True when the wait is over. Airflow calls poke repeatedly. The arguments every sensor shares:

Argument Default (Airflow 3.3) Meaning
poke_interval 60 seconds Time between checks.
timeout 7 days ([sensors] default_timeout, 604800 s) Total time allowed for the wait, across all pokes and reschedules. Then the sensor fails with AirflowSensorTimeout.
mode "poke" "poke" or "reschedule"; see the next section.
soft_fail False On timeout, mark the sensor skipped instead of failed.
exponential_backoff False Grow the interval between pokes.
max_wait none Cap for the interval when backing off.
deferrable False on most sensors ([operators] default_deferrable) Hand the wait to the triggerer.

timeout is not the same as execution_timeout. timeout limits the whole wait; execution_timeout limits one attempt, and a sensor that hits execution_timeout can be retried, while one that hits timeout fails without retrying.

The quickest custom sensor is @task.sensor (or PythonSensor). Return a PokeReturnValue to also push a value to XCom when the wait ends:

from airflow.sdk import DAG, dag, task, PokeReturnValue
from airflow.providers.standard.sensors.python import PythonSensor

poke_count = {"n": 0}

def batch_is_ready() -> bool:
    poke_count["n"] += 1
    return poke_count["n"] >= 3          # becomes true on the third check

@dag(schedule=None, start_date=datetime(2026, 1, 1))
def sensor_basics():
    wait_py = PythonSensor(task_id="wait_for_flag", python_callable=batch_is_ready,
                           poke_interval=1, timeout=30)

    @task.sensor(poke_interval=1, timeout=30, mode="reschedule")
    def wait_for_batch() -> PokeReturnValue:
        batch = {"batch_id": 42}                      # pretend we looked this up
        return PokeReturnValue(is_done=True, xcom_value=batch)

    @task
    def process(batch: dict) -> str:
        return f"processing batch {batch['batch_id']}"

    wait_py >> process(wait_for_batch())

print(run_dag(sensor_basics())[1:])
print("pokes needed:", poke_count["n"])
({'process': 'success', 'wait_for_batch': 'success', 'wait_for_flag': 'success'}, {'process': 'processing batch 42', 'wait_for_batch': {'batch_id': 42}})
pokes needed: 3

Pitfalls

  • Leaving timeout at the default 7 days. A missing upstream file then holds a slot (in poke mode) for a week before anyone hears about it. Set a timeout that matches how late the data can reasonably be, and alert when it passes.
  • poke_interval of a few seconds on an external API: thousands of calls a day per sensor. Use backoff or deferral.
  • Side effects inside poke. It runs many times, so it must be a read-only check.
  • A wall of sensors at the start of every DAG run, all waiting on the same thing. One sensor (or an asset-based schedule) is enough.

In interviews

“What is a sensor and what can go wrong with them?” A good answer explains poke, the interval and timeout, and then the real issue: sensors occupy worker slots while waiting unless you use reschedule or deferrable mode. Knowing that soft_fail turns a timeout into a skip, and the difference between timeout and execution_timeout, rounds it off.

Poke vs reschedule mode

What it is

The mode decides what a sensor does between checks:

  • poke (default): the task stays running on its worker and sleeps between pokes. It holds a worker slot (and a pool slot) for the whole wait.
  • reschedule: after a failed poke the task gives up its slot, goes to the up_for_reschedule state, and the scheduler starts it again after poke_interval. Between pokes it uses no worker.

How it works

Poke Reschedule Deferrable (next sections)
Worker slot while waiting held the whole time released between pokes released; the triggerer waits
Cost per check a function call a full task start (process, DAG parse, context) an async coroutine on the triggerer
Good for short waits, intervals of seconds long waits with intervals of minutes long waits at scale, when a deferrable version exists
State while waiting running up_for_reschedule deferred

Rule of thumb: if the wait is shorter than a few minutes, poke; if it is longer, and especially with a poke_interval of 60 seconds or more, reschedule or (better) deferrable. Reschedule with a very short interval is wasteful, because each check is a full task start.

from airflow.providers.standard.sensors.filesystem import FileSensor

with DAG("modes_demo", schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False) as modes_dag:
    FileSensor(task_id="poke_mode", fs_conn_id="fs_landing", filepath="orders_{{ ds }}.csv",
               mode="poke", poke_interval=1, timeout=30)
    FileSensor(task_id="reschedule_mode", fs_conn_id="fs_landing", filepath="orders_{{ ds }}.csv",
               mode="reschedule", poke_interval=1, timeout=30)
    FileSensor(task_id="late_but_optional", fs_conn_id="fs_landing", filepath="returns_{{ ds }}.csv",
               mode="reschedule", poke_interval=1, timeout=3, soft_fail=True)
    FileSensor(task_id="late_and_required", fs_conn_id="fs_landing", filepath="returns_{{ ds }}.csv",
               mode="reschedule", poke_interval=1, timeout=3)

print(run_dag(modes_dag, logical_date=MARCH_1)[:2])
('failed', {'late_and_required': 'failed', 'late_but_optional': 'skipped', 'poke_mode': 'success', 'reschedule_mode': 'success'})

Both modes found the file. For the missing file, the 3-second timeout fired: with soft_fail=True the sensor was skipped, without it the sensor failed and so did the DAG run.

Pitfalls

  • Sensor deadlock: with poke mode, enough sensors waiting on work that needs a slot can fill every slot, and the work they wait for can never start. Reschedule or deferrable mode, or a dedicated pool for sensors, prevents it.
  • Reschedule mode with poke_interval of a few seconds, which costs more than poke mode.
  • Expecting reschedule mode to reset the timeout. The timeout counts from the first poke across all reschedules.
  • Reschedule mode is not available for every operator; it is a sensor feature.

In interviews

“Poke vs reschedule?” is one of the most asked Airflow questions. Explain slot usage, the cost of each check, the up_for_reschedule state and the rule of thumb, then add that deferrable operators are the modern answer for long waits because a single triggerer process can wait on many tasks at once.

ExternalTaskSensor

What it is

ExternalTaskSensor waits for a task (or a whole DAG run, or a task group) in another DAG to reach a state, by default success. It is the classic way to make DAG B wait for DAG A when they are owned by different teams or run on different schedules.

How it works

  • Import from the standard provider: airflow.providers.standard.sensors.external_task (Airflow 2: airflow.sensors.external_task).
  • external_dag_id plus external_task_id, external_task_ids or external_task_group_id. With no task id it waits for the DAG run itself.
  • Date alignment: by default it looks for the external run with the same logical date as the current run. If the schedules differ, shift with execution_delta (a fixed timedelta) or execution_date_fn (a function returning the date(s) to check). Both parameters keep their Airflow 2 names in the standard provider even though the concept is now the logical date.
  • allowed_states, failed_states and skipped_states control which external states mean success, failure (stop waiting early) or skip.
  • check_existence=True fails fast if the external DAG or task does not exist, instead of waiting until the timeout.
  • deferrable=True waits in the triggerer.
  • ExternalTaskMarker is the companion for the other direction: when you clear a task in DAG A with “downstream” selected, the marker also clears the dependent tasks in DAG B.
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.providers.standard.sensors.external_task import ExternalTaskSensor

with DAG("upstream_etl", schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False) as upstream:
    EmptyOperator(task_id="publish")

with DAG("downstream_report", schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False) as downstream:
    wait = ExternalTaskSensor(
        task_id="wait_for_upstream",
        external_dag_id="upstream_etl",
        external_task_id="publish",
        failed_states=["failed"],          # stop waiting at once if upstream failed
        mode="reschedule", poke_interval=1, timeout=5,
    )
    wait >> EmptyOperator(task_id="build_report")

print("upstream 1 Mar:  ", run_dag(upstream, logical_date=MARCH_1)[:2])
print("downstream 1 Mar:", run_dag(downstream, logical_date=MARCH_1)[:2])
print("downstream 2 Mar:", run_dag(downstream, logical_date=MARCH_1.add(days=1))[:2])
upstream 1 Mar:   ('success', {'publish': 'success'})
downstream 1 Mar: ('success', {'build_report': 'success', 'wait_for_upstream': 'success'})
downstream 2 Mar: ('failed', {'build_report': 'upstream_failed', 'wait_for_upstream': 'failed'})

For 2 March there was no upstream run with that logical date, so the sensor timed out and the report never ran.

Pitfalls

  • Mismatched schedules. An hourly DAG waiting on a daily DAG, or two DAGs with different start_date times, look for logical dates that never exist and wait until the timeout. Use execution_delta or execution_date_fn.
  • In Airflow 3 a cron schedule creates runs whose logical date is the run time (no data interval by default). If one DAG uses data intervals and the other does not, the logical dates differ by one interval. Check both before relying on the default alignment.
  • Manual runs of the downstream DAG have a logical date of “now”, which matches no upstream run.
  • Chains of sensors across many DAGs become a hidden dependency graph that nobody can see in one place. Consider assets (data-aware scheduling) instead: the producer declares what it updates and the consumer is scheduled on it, with no polling. The scheduling lesson covers assets.

In interviews

“How do you create a dependency between two DAGs?” Give three options and their trade-offs: ExternalTaskSensor (pull, polling, needs aligned dates), TriggerDagRunOperator (push from the upstream DAG), and assets (data-aware scheduling, the modern default). Mention date alignment as the classic ExternalTaskSensor bug and ExternalTaskMarker for clearing across DAGs.

FileSensor

What it is

FileSensor waits for a file or folder to exist on a filesystem the worker can see (a local disk or a mounted network share). For object stores you use the provider sensors instead, such as S3KeySensor (Amazon), GCSObjectExistenceSensor (Google) or WasbBlobSensor (Azure), which follow the same pattern.

How it works

  • Import: airflow.providers.standard.sensors.filesystem.FileSensor (Airflow 2: airflow.sensors.filesystem).
  • fs_conn_id (default fs_default) names a connection of type fs whose extra holds a base path; filepath is joined to it. filepath is templated, so orders_{{ ds }}.csv waits for that run’s file.
  • The path is matched with glob, so wildcards work (orders_{{ ds }}_*.csv), and recursive=True allows **. A matching directory counts only if it contains files.
  • deferrable=True waits in the triggerer using FileTrigger.
with DAG("file_sensor_demo", schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False) as file_dag:
    FileSensor(task_id="exact_name", fs_conn_id="fs_landing", filepath="orders_{{ ds }}.csv",
               poke_interval=1, timeout=10)
    FileSensor(task_id="wildcard", fs_conn_id="fs_landing", filepath="orders_*.csv",
               poke_interval=1, timeout=10)
    FileSensor(task_id="deferred", fs_conn_id="fs_landing", filepath="orders_{{ ds }}.csv",
               deferrable=True, poke_interval=1)

print(run_dag(file_dag, logical_date=MARCH_1)[:2])
('success', {'deferred': 'success', 'exact_name': 'success', 'wildcard': 'success'})

Pitfalls

  • The file exists but is not complete yet. A sensor fires as soon as the name appears, possibly while the writer is still copying it. Have producers write to a temporary name and rename at the end, or wait for a marker file such as _SUCCESS.
  • The path exists on the machine where you tested, not on the workers. With Celery or Kubernetes workers each one needs the mount.
  • Wildcards that match yesterday’s leftover file. Template the date into the pattern.
  • Using FileSensor for S3 or GCS paths. Use the provider sensor, which talks to the object store API.

In interviews

“How do you start processing when an upstream file lands?” Mention a file or object-store sensor with a templated path, a sensible timeout, reschedule or deferrable mode, and a completion marker so you do not read half-written files. For event-driven designs, also mention asset-based scheduling or an event that triggers the DAG through the REST API.

Deferrable operators and triggers

What it is

A deferrable operator can pause itself mid-task and hand the waiting to a trigger: a small asynchronous Python object that runs in a separate triggerer process. While deferred, the task holds no worker slot. When the trigger fires an event, the scheduler resumes the task on a worker, which continues in a callback method. One triggerer can run thousands of triggers concurrently because they are asyncio coroutines.

How it works

  1. The operator’s execute starts the work (for example submits a Spark job) and calls self.defer(trigger=..., method_name="execute_complete").
  2. The task instance goes to the deferred state and is removed from the worker. The trigger is serialised (serialize() returns its import path and kwargs) and stored in the database.
  3. The triggerer imports the trigger class and runs its async def run(), which yields a TriggerEvent when the condition is met.
  4. The task is scheduled again and Airflow calls execute_complete(context, event) on a worker. Its return value is the task’s result.

Requirements: run at least one airflow triggerer process (it is part of a standard deployment), and make the trigger class importable on the triggerer, so define it in a package, not in the DAG file. Many provider operators and sensors already support deferrable=True (for example TimeDeltaSensor, DateTimeSensorAsync, FileSensor, ExternalTaskSensor, and most job-running operators in the Amazon, Google and Databricks providers). Set [operators] default_deferrable = True to make deferrable the default for operators that support it.

A minimal deferrable operator and trigger:

import sys, types
from airflow.sdk import BaseOperator
from airflow.triggers.base import BaseTrigger, TriggerEvent

# The triggerer re-imports a trigger by its dotted path, so it must live in an importable module.
# In a project that is a file in your package (for example my_company/triggers.py); here an
# in-memory module stands in for it.
job_triggers = types.ModuleType("job_triggers")
sys.modules["job_triggers"] = job_triggers

class JobDoneTrigger(BaseTrigger):
    """Polls a (pretend) job API asynchronously until the job finishes."""
    def __init__(self, job_id: str, polls_needed: int = 3, interval: float = 0.5):
        super().__init__()
        self.job_id, self.polls_needed, self.interval = job_id, polls_needed, interval

    def serialize(self):
        # import path + kwargs: how the triggerer rebuilds this object
        return (f"{type(self).__module__}.{type(self).__qualname__}",
                {"job_id": self.job_id, "polls_needed": self.polls_needed, "interval": self.interval})

    async def run(self):
        for poll in range(1, self.polls_needed + 1):
            await asyncio.sleep(self.interval)   # never use blocking calls such as time.sleep here
        yield TriggerEvent({"job_id": self.job_id, "status": "SUCCEEDED", "polls": poll})

JobDoneTrigger.__module__ = "job_triggers"
job_triggers.JobDoneTrigger = JobDoneTrigger

class RunJobOperator(BaseOperator):
    def __init__(self, *, job_name: str, **kwargs):
        super().__init__(**kwargs)
        self.job_name = job_name

    def execute(self, context):
        job_id = f"{self.job_name}-123"                 # pretend we submitted the job
        self.defer(trigger=JobDoneTrigger(job_id=job_id), method_name="execute_complete")

    def execute_complete(self, context, event=None):
        if event["status"] != "SUCCEEDED":
            raise RuntimeError(f"job failed: {event}")
        return event

with DAG("deferrable_demo", schedule=None, start_date=datetime(2026, 1, 1)) as defer_dag:
    RunJobOperator(task_id="spark_job", job_name="orders")

print(run_dag(defer_dag)[1:])
({'spark_job': 'success'}, {'spark_job': {'job_id': 'orders-123', 'status': 'SUCCEEDED', 'polls': 3}})

Built-in triggers you will reuse: DateTimeTrigger and TimeDeltaTrigger (wait for a time), FileTrigger (wait for a file), and the provider triggers that poll job APIs. defer() also accepts a timeout, after which the task fails.

Pitfalls

  • Blocking code in run(). A time.sleep or a synchronous HTTP call blocks the triggerer’s event loop and stalls every other trigger on it. Use async libraries or asyncio.to_thread. The triggerer reports blocked-loop warnings for this reason.
  • Trigger classes defined in the DAG file. The triggerer imports by path and cannot see your DAG’s module namespace.
  • No triggerer running. Deferred tasks then sit in the deferred state forever.
  • Keeping state on self between execute and execute_complete. They run in different processes; pass what you need through the trigger’s event or defer(kwargs=...).
  • Triggers can be run more than once (for example after a triggerer restart), so run() must be safe to repeat.

In interviews

Explain the lifecycle in four steps (execute, defer, trigger runs on the triggerer, resume in a callback) and why it scales: async coroutines are cheap, worker slots are not. Interviewers may ask how you would make a custom operator deferrable, what the triggerer is, and what happens if it crashes (triggers are reassigned and re-run, which is why they must be idempotent).

Smart Sensors (removed)

What it was

Smart Sensors were an early attempt, added in Airflow 2.0, to stop sensors wasting worker slots. Sensors that opted in did not poke themselves. Instead they registered their poke arguments in the database, and a small number of special “smart sensor” DAG tasks ran the pokes for many sensors in batches.

Why it went

Smart Sensors were deprecated in Airflow 2.2 and removed in Airflow 2.4. The design only worked for sensors whose arguments could be stored and re-executed by a central task, it was hard to operate, and the same problem was solved more generally by deferrable operators and triggers (added in 2.2), which work for sensors and for any long-running operator, such as waiting for a Spark job. There is no automatic migration path: you switch to the deferrable version of the sensor (deferrable=True or an ...Async sensor), use reschedule mode, or accept poke mode for short waits.

In interviews

You may still be asked about Smart Sensors because older blog posts and certification material mention them. A strong answer: “They batched sensor pokes into a few central tasks to save slots; they were removed in 2.4 because deferrable operators with the triggerer solve the same problem more generally and are now the recommended approach.”

Practice questions

A DAG has 50 sensors in poke mode waiting on files that arrive hours later, and other DAGs stop running. Why, and how do you fix it?

Each poke-mode sensor holds a worker (and pool) slot for the whole wait, so the sensors use up capacity that other tasks need, possibly including the tasks that would produce the files (a sensor deadlock). Switch them to deferrable mode (best) or reschedule mode, give them a sensible timeout, and optionally put sensors in their own pool to cap how many slots they can take.

What is the difference between timeout and execution_timeout on a sensor?

timeout bounds the total waiting time across all pokes and reschedules; when it passes, the sensor fails with AirflowSensorTimeout and is not retried (or is skipped with soft_fail=True). execution_timeout bounds a single attempt of the task; hitting it fails that try, and normal retries apply.

An ExternalTaskSensor waits forever even though the upstream DAG succeeded. What do you check?

Date alignment. The sensor looks for an upstream run with the same logical date by default. Different schedules, different start times, or one DAG using data intervals and the other not, mean no run matches. Fix with execution_delta or execution_date_fn, add check_existence=True, and set failed_states and a timeout so it does not wait silently. Consider assets instead.

How does a deferrable operator free its worker slot?

In execute it calls self.defer(trigger=..., method_name=...). Airflow stores the serialised trigger, sets the task to deferred and ends the worker process. The triggerer runs the trigger’s async run() method; when it yields an event, the task is scheduled again and the named method runs on a worker with the event.

Why must code in a trigger’s run() be asynchronous?

All triggers on a triggerer share one asyncio event loop. A blocking call stops the loop, so every other trigger stalls until it returns. Use async clients, await asyncio.sleep, or push blocking work to a thread with asyncio.to_thread.

What were Smart Sensors and what replaced them?

An Airflow 2.0 feature that batched the pokes of many sensors into a few central smart-sensor tasks to save worker slots. They were deprecated in 2.2 and removed in 2.4, replaced by deferrable operators and the triggerer, which solve the same problem for both sensors and long-running operators.

A FileSensor succeeds but the next task reads a truncated file. What happened?

The sensor fired as soon as the file name existed, while the producer was still writing. Make the producer write to a temporary name and rename atomically when done, or wait for a completion marker such as _SUCCESS, and point the sensor at the final name or marker.

Key takeaways

  • A sensor repeatedly checks a condition; set poke_interval and a realistic timeout (the default is 7 days) and use soft_fail when the wait is optional.
  • Poke mode holds a worker slot for the whole wait; reschedule mode releases it between checks; deferrable mode hands the wait to the triggerer and scales best.
  • ExternalTaskSensor matches runs by logical date, so align schedules with execution_delta or execution_date_fn, or use assets instead.
  • FileSensor reads a fs connection’s base path plus a templated, glob-able filepath; guard against half-written files with markers or atomic renames.
  • Deferrable operators call defer(), a trigger waits asynchronously in the triggerer, and the task resumes in a callback; keep trigger code non-blocking and importable.
  • Smart Sensors were removed in Airflow 2.4; deferrable operators are their replacement.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Examples run on Apache Airflow 3.3.2 with apache-airflow-providers-standard 1.20.0, Python 3.11 and a SQLite metadata database; dag.test() runs triggers in-process, so no separate triggerer was needed. Smart Sensors are described from the Airflow 2.2 documentation and cannot run on any current version.

Progress is saved in this browser only. No account needed.

Search
Filter by type