How do you get alerted when an Airflow task fails, and what do the callback arguments give you?
answer
- one path sends mail, one runs your code
- a callable per lifecycle point
- it receives the run's context dict
- the log URL lives on the task instance
- it only fires after the last attempt
basics
~20 sAirflow fires per-task callables such as on_failure_callback, on_retry_callback and on_success_callback, plus email_on_failure to a configured SMTP address. Callbacks receive a context dictionary with the DAG, task, run and log URL, which you forward to Slack, PagerDuty or a webhook.
solid answer
~40 sAirflow gives you two alerting surfaces. The built-in one is `email_on_failure` / `email_on_retry` with the `email` argument, which needs SMTP configured and only sends mail. The one teams actually use is **callbacks**: `on_failure_callback`, `on_retry_callback`, `on_success_callback` and `on_execute_callback` on any operator, plus `on_failure_callback` at the DAG level for run-level events. Each is a callable receiving the task **context** dict — `dag`, `task`, `ti`, `run_id`, `logical_date`, the exception, and the log URL — from which you build a Slack or PagerDuty payload. Set them once in `default_args` so every task inherits, rather than per task. Two judgment points matter: `on_failure_callback` fires only after the *last* retry, so retry counts add latency to your page; and the callback runs in the worker, so a slow or throwing callback affects the task's own teardown.
code
python · 33 linesfrom datetime import datetime, timedelta
from airflow.decorators import dag, task
def notify_failure(context):
ti = context["ti"]
msg = (
f"FAILED {ti.dag_id}.{ti.task_id}\n"
f"run={context['run_id']} try={ti.try_number}\n"
f"error={context.get('exception')}\n"
f"{ti.log_url}"
)
try:
post_to_slack(msg)
except Exception:
pass # never let the notifier mask the real failure
def note_retry(context):
ti = context["ti"]
increment_counter("airflow.retry", tags=[ti.dag_id, ti.task_id])
@dag(
schedule="@daily",
start_date=datetime(2024, 1, 1),
catchup=False,
on_failure_callback=notify_failure, # fires once for the DAG run
default_args={
"retries": 2,
"retry_delay": timedelta(minutes=5),
"on_retry_callback": note_retry, # metric, not a page
},
)
def daily_sales():
...go deeper
Recall the argument names — email_on_failure, on_failure_callback, on_retry_callback, on_success_callback — and that a callback is just a Python function receiving the run's context.
Explain what the context dict carries, how to build a useful message from ti.log_url and the exception, and why you wire the callback once in default_args rather than per task.
Show the timing judgment: failure callbacks fire only after the last retry, DAG-level beats per-task for paging, and failure alerting cannot see a run that never happened.
Own the alerting policy across many DAGs — what pages versus what goes to a channel, freshness checks alongside failure alerts, and the metrics pipeline (StatsD/OpenTelemetry) that answers questions no individual notification can.
## The two mechanisms Airflow ships with a plain email path and a general callback path. **Email** is controlled by three `BaseOperator` arguments: `email` (address or list), `email_on_failure` and `email_on_retry`. It requires a working SMTP configuration in `airflow.cfg` or the corresponding environment variables. It is the oldest mechanism, it only sends mail, and it has no routing, deduplication or escalation. It is fine for a hobby deployment and rarely sufficient for a production one. **Callbacks** are arbitrary Python callables invoked at defined lifecycle points: - `on_execute_callback` — just before the task's execute begins. - `on_success_callback` — after a successful run. - `on_retry_callback` — on each retry, i.e. every time the instance enters `up_for_retry`. - `on_failure_callback` — when the instance reaches terminal `failed`, after retries are exhausted. The same `on_failure_callback` and `on_success_callback` names also exist on the DAG object, where they fire for the **DAG run** rather than an individual task — useful for a single 'pipeline X failed' notification instead of one message per failed task. ## The context dictionary Every callback receives one argument, conventionally called `context`: a dict of the same values Jinja templates see. The useful keys for alerting are `dag`, `task`, `ti` (the task instance, from which you get `try_number`, `state` and `log_url`), `run_id`, `logical_date`, `data_interval_start` / `data_interval_end`, and `exception` — the error object that caused the failure. A good alert names the DAG and task, the run's interval, the attempt number, a one-line error summary, and a deep link built from `ti.log_url` so the responder lands directly on the failing log. ```python def notify_failure(context): ti = context["ti"] send_slack( f"FAILED {ti.dag_id}.{ti.task_id} " f"run={context['run_id']} try={ti.try_number}\n" f"{context.get('exception')}\n{ti.log_url}" ) ``` Wiring it once in `default_args` is the right default so nobody ships a silently-failing task: ```python default_args = {"on_failure_callback": notify_failure, "retries": 2} ``` For Slack specifically, the community provider ships helpers (a Slack webhook notifier and the Slack operators) so you rarely have to hand-roll the HTTP call, but the shape is the same: a callable that closes over a connection id and formats the context. ## Timing: retries delay the page The most important operational detail is that `on_failure_callback` fires only when the task reaches terminal `failed` — that is, after every retry has been exhausted. A task with `retries=4` and `retry_delay=timedelta(minutes=10)` will not alert for something like forty minutes after the first symptom. If that is too slow for a pipeline with a downstream commitment, you have three levers: reduce retries for that task, emit a lower-severity signal from `on_retry_callback` so you can *see* the flapping before it pages, or attach the real deadline to an SLA rather than to failure. Using `on_retry_callback` as a metric rather than a page is a good pattern generally: increment a counter or post to a low-noise channel, so a task quietly succeeding on its third attempt every night becomes visible instead of invisible. ## Where callbacks run, and why it matters A task-level callback executes in the task's own process on the worker, as part of finishing the instance. Three consequences follow. First, the callback must be **fast** — a slow HTTP call to a paging provider holds a worker slot. Second, it must not **throw**: an exception inside a failure callback is logged but leaves you with a failed task and no notification, the worst combination, so wrap the delivery in a try/except. Third, it needs whatever credentials it uses to be reachable from the worker, normally via an Airflow connection rather than an inline secret. A callback also cannot fire if the process never got to run it — the zombie case, where the worker was OOM-killed or the pod evicted. There the scheduler fails the instance and the failure callback does get invoked from that path, but any assumption that your notifier always runs inside the dying process is wrong. ## What good alerting looks like Beyond the mechanics, interviewers want the judgment layer: - **Route by severity, not by task.** Everything paging is the same as nothing paging. A failed reporting refresh belongs in a channel; a failed regulatory extract pages. - **Alert on the pipeline, not each task.** A DAG-level `on_failure_callback` avoids twenty messages when a fan-out of twenty mapped tasks all fail on the same upstream outage. - **Alert on staleness too.** A DAG that never started — because it was paused, or the scheduler was down — produces no failure and therefore no callback. Failure alerting is blind to absence; freshness checks on the output table or an SLA are what catch it. - **Include the deep link.** The single biggest quality difference between alerts is whether the responder can reach the log in one click. - **Emit metrics as well as messages.** Airflow can emit StatsD metrics and, in recent versions, OpenTelemetry, giving you dashboards for task duration and failure counts that no per-event notification provides.
- Why might a failure alert arrive forty minutes after the pipeline actually broke?Because `on_failure_callback` fires only on terminal failure, after all retries are exhausted. With `retries=4` and a ten-minute `retry_delay`, roughly forty minutes elapse before the instance reaches `failed`. If that is unacceptable, cut retries for that task, emit a low-severity signal from `on_retry_callback`, or attach the deadline to an SLA instead.
- What is the difference between a task-level and a DAG-level on_failure_callback?The task-level one fires per failing task instance; the DAG-level one fires for the DAG run. With a fan-out of twenty mapped tasks all failing on the same upstream outage, task-level callbacks produce twenty messages while the DAG-level one produces a single 'this pipeline failed' notification. Most teams use the DAG level for paging and the task level for detail.
- Your DAG is paused by mistake and produces no runs for two days. Does failure alerting catch it?No. No run means no task, no failure, and no callback — failure alerting is blind to absence. Catching it needs something that asserts presence: a freshness or completion check on the output table, an SLA/deadline mechanism, or an external monitor watching for the expected run. This is the classic gap that makes 'we alert on failures' insufficient.
- What should you be careful about inside the callback function itself?It runs in the worker as part of finishing the task, so keep it fast — a slow paging API call holds a worker slot — and wrap the delivery in try/except, because an exception inside the failure callback leaves you with a failed task and no notification at all. Pull credentials from an Airflow connection rather than inlining them.
saying these in an interview costs you the question
- Relies on email_on_failure alone with no SMTP or routing
- Thinks on_failure_callback fires on the first failed attempt
- Assumes failure alerts catch a DAG that never ran
- Puts a slow or unguarded HTTP call inside the callback
- Pages on every task failure with no severity routing