In Airflow, what problem do deferrable operators and the triggerer solve?
answer
- waiting should not cost a worker
- something else holds the wait for you
- one process, many coroutines
- the task's state while it waits is not running
- self.defer, and the component that must be running
basics
~20 sThey remove the cost of waiting. A deferrable operator starts its work, hands a small async trigger to Airflow's triggerer process and releases its worker slot; the task resumes only when the trigger fires, so thousands of waits cost one process instead of thousands of slots.
solid answer
~50 sA task that waits — for a file, for a remote Spark or BigQuery job, for another DAG — normally burns a worker slot for the whole wait. A **deferrable** operator instead calls `self.defer(trigger=..., method_name=...)`, which raises `TaskDeferred`: the worker process exits, the task instance goes to the `deferred` state, and the trigger is handed to the **triggerer**, a separate Airflow component running many triggers concurrently in one asyncio event loop. When the trigger yields a `TriggerEvent`, the scheduler re-queues the task and it resumes on a worker at `method_name`. Many provider operators and sensors accept `deferrable=True` for exactly this. Compared with `mode='reschedule'`, which frees the slot but still pays a full task startup per poke, deferral pays neither — and it covers long-running remote jobs, not just pokeable conditions. The requirements: you must actually run `airflow triggerer`, and trigger code must be async and non-blocking, because one blocking call stalls every other trigger in that loop.
code
python · 16 linesfrom airflow.models import BaseOperator
from airflow.triggers.temporal import DateTimeTrigger
class WaitUntil(BaseOperator):
def __init__(self, *, moment, **kwargs):
super().__init__(**kwargs)
self.moment = moment
def execute(self, context):
# returns nothing: raises TaskDeferred and frees the worker slot
self.defer(trigger=DateTimeTrigger(moment=self.moment), method_name="resume")
def resume(self, context, event=None):
self.log.info("woken by %s", event)
return "done"go deeper
Recall that passing deferrable=True lets a waiting task release its worker while a separate Airflow component holds the wait, and that the component is called the triggerer.
Explain the mechanics: execute() calls self.defer with a trigger, the task goes to the deferred state, the triggerer awaits it in an asyncio loop, and the task resumes on a worker at the named method.
Be ready to justify deferral against reschedule mode with numbers of concurrent waits, and to name the operational requirements: run the triggerer, keep triggers non-blocking and serialisable, monitor triggerer health.
Own the capacity story: how much of your fleet exists only to wait, whether deferral or event-driven triggering removes it, and what deploying and making the triggerer highly available costs.
## The problem Waiting is the dominant activity in most orchestration deployments. A DAG waits for a vendor file, submits a warehouse query and waits for it, launches a Spark job and waits for it, waits for an upstream DAG. In the classic model, each of those waits is a running task instance holding a worker slot. Scale that to a few hundred concurrent pipelines and your worker fleet is sized not by the work you do but by the waiting you do — and when capacity runs out, real tasks queue behind sleeping ones. `mode='reschedule'` fixed half of this for sensors: give the slot back between checks. It does not help a task that has *submitted* something and must stay attached to it, and it still pays a complete task launch — scheduler pickup, process fork, DAG parse, connection setup — on every single check. ## What deferral does Deferrable operators, introduced in Airflow 2.2 along with the **triggerer** component, split such a task into two phases with an asynchronous wait in the middle. 1. `execute(context)` does the synchronous part — submit the job, record the id — and then calls `self.defer(trigger=SomeTrigger(...), method_name="resume")`. 2. `defer()` raises `TaskDeferred`. The worker process ends and the task instance is marked `deferred`. **No worker slot is held.** 3. The trigger object — small, serialisable, described by a classpath plus kwargs — is picked up by the triggerer. The triggerer runs an asyncio event loop in which many triggers coexist; each is an async generator that awaits its condition and finally yields a `TriggerEvent`. 4. On that event the scheduler moves the task back to a queued state; a worker picks it up and calls `resume(context, event)`, which finishes the task. Because the waits are coroutines rather than processes, one triggerer can hold a very large number of them concurrently. The unit cost of a wait drops from "a slot" to "an entry in an event loop". ```python from airflow.models import BaseOperator from airflow.triggers.temporal import DateTimeTrigger class WaitUntil(BaseOperator): def __init__(self, *, moment, **kwargs): super().__init__(**kwargs) self.moment = moment def execute(self, context): self.defer(trigger=DateTimeTrigger(moment=self.moment), method_name="resume") def resume(self, context, event=None): self.log.info("resumed at %s", event) ``` Most of the time you write none of this: you pass `deferrable=True` to a provider operator or sensor that already supports it, and the async variants such as `DateTimeSensorAsync` and `TimeDeltaSensorAsync` exist for the built-in waits. ## What you must get right - **Run the triggerer.** `airflow triggerer` is a component alongside the scheduler, webserver and workers. Deferring with no triggerer running means tasks sit in `deferred` forever. This is the single most common self-inflicted failure with the feature. - **Triggers must not block.** They share one event loop. A synchronous HTTP call or a `time.sleep()` inside a trigger stalls every other deferred task on that process. Use async clients; keep the trigger tiny. - **Triggers must be serialisable.** A trigger is stored as a classpath plus keyword arguments so it can be reconstructed; you cannot stash an open client or a live connection in it. State that must survive the wait goes in the trigger's arguments. - **`execute()` may run twice in effect.** The pre-defer phase runs, the task ends, and the resume method runs later in a *different* process — nothing in memory survives. Anything the resume needs must be passed via the trigger or re-fetched. - **Triggerer capacity is real.** It is highly concurrent, not infinite, and it is another component to deploy, monitor and make highly available. ## How to compare the three waiting strategies in an interview - **Poke mode**: simplest, holds a slot the whole time, fine for waits of seconds. - **Reschedule mode**: frees the slot between checks, costs a full task launch per check, sensors only, good for hour-scale waits when deferral is unavailable. - **Deferrable**: frees the slot and pays no per-check task launch, works for submitted remote jobs as well as pokeable conditions, requires the triggerer and async-safe trigger code. And the strongest answer above all three: if the event has a producer you control, arrange for the producer to trigger the consumer so nothing waits at all.
- How does a deferrable sensor differ from the same sensor in reschedule mode?Reschedule mode ends the task after each failed poke and pays a full task launch — scheduler pickup, process start, DAG parse — for the next check. Deferral hands one coroutine to the triggerer and pays nothing until the condition fires. Deferral also covers waits on submitted remote jobs, which reschedule mode cannot express.
- What happens if a trigger performs a blocking network call?It stalls the shared asyncio event loop, so every other deferred task on that triggerer stops progressing until the call returns. Triggers must use async clients and stay small; anything heavy belongs in the operator's execute or resume phase, which runs on a worker.
- Where does state live across the deferral boundary?Only in the trigger's serialised arguments and in whatever the task persisted before deferring, such as an XCom or a remote job id. The worker process that called defer is gone, so no in-memory object survives; the resume method runs fresh in a different process and must reconstruct what it needs.
saying these in an interview costs you the question
- Enables deferrable operators without running the triggerer component
- Puts a blocking HTTP call or sleep inside a trigger
- Assumes deferrable is just a faster reschedule mode
- Expects in-memory state from execute to survive into the resume method
- Thinks the triggerer executes the task's actual work