Airflow courseLesson 3 of 5
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.
On this page
- Running the examples
- Sensor basics
- What it is
- How it works
- Pitfalls
- In interviews
- Poke vs reschedule mode
- What it is
- How it works
- Pitfalls
- In interviews
- ExternalTaskSensor
- What it is
- How it works
- Pitfalls
- In interviews
- FileSensor
- What it is
- How it works
- Pitfalls
- In interviews
- Deferrable operators and triggers
- What it is
- How it works
- Pitfalls
- In interviews
- Smart Sensors (removed)
- What it was
- Why it went
- In interviews
- Practice questions
- 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
timeoutat 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_intervalof 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_intervalof 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_idplusexternal_task_id,external_task_idsorexternal_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) orexecution_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_statesandskipped_statescontrol which external states mean success, failure (stop waiting early) or skip.check_existence=Truefails fast if the external DAG or task does not exist, instead of waiting until the timeout.deferrable=Truewaits 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_datetimes, look for logical dates that never exist and wait until the timeout. Useexecution_deltaorexecution_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(defaultfs_default) names a connection of typefswhoseextraholds a basepath;filepathis joined to it.filepathis templated, soorders_{{ ds }}.csvwaits for that run’s file.- The path is matched with
glob, so wildcards work (orders_{{ ds }}_*.csv), andrecursive=Trueallows**. A matching directory counts only if it contains files. deferrable=Truewaits in the triggerer usingFileTrigger.
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
FileSensorfor 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
- The operator’s
executestarts the work (for example submits a Spark job) and callsself.defer(trigger=..., method_name="execute_complete"). - 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. - The triggerer imports the trigger class and runs its
async def run(), which yields aTriggerEventwhen the condition is met. - 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(). Atime.sleepor a synchronous HTTP call blocks the triggerer’s event loop and stalls every other trigger on it. Use async libraries orasyncio.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
selfbetweenexecuteandexecute_complete. They run in different processes; pass what you need through the trigger’s event ordefer(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_intervaland a realistictimeout(the default is 7 days) and usesoft_failwhen 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_deltaorexecution_date_fn, or use assets instead. - FileSensor reads a
fsconnection’s base path plus a templated, glob-ablefilepath; 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.
Progress is saved in this browser only. No account needed.