skip to content

Monitoring and Retries

Retries with backoff, SLA misses, callbacks, and the task states you read when something is stuck. Interviewers ask what a zombie task is and how you alert on failures, because pipelines fail routinely and silence is the real bug.

on this pageshow

questions

5

In Airflow, what do the task-level retries and retry_delay arguments control?

level: juniorimportance: must knowfreq 78%

answer

  1. one number, one duration
  2. how many more tries, and how long between
  3. the state a task sits in while waiting
  4. exponential variant exists for hammering services
  5. the count excludes the first attempt

basics

~10 s

In Airflow, retries sets how many extra attempts a task instance gets after its first failure, and retry_delay is the wait between attempts. Between attempts the task sits in the up_for_retry state.

solid answer

~40 s

In Airflow, `retries` is the number of *additional* attempts a task instance gets after the initial run fails, so `retries=2` means up to three executions in total. `retry_delay` is a `timedelta` controlling how long the scheduler waits before the next attempt; while waiting, the task instance sits in the `up_for_retry` state and the run stays in progress rather than failing. Both can be set per task or applied to every task in the DAG through `default_args`. Adding `retry_exponential_backoff=True` makes each successive delay grow instead of staying constant, and `max_retry_delay` caps that growth. Only when the last retry fails does the task instance move to `failed`, at which point downstream tasks are skipped or marked upstream_failed depending on their `trigger_rule`.

code

python · 26 lines
python
from datetime import datetime, timedelta
from airflow.decorators import dag, task

@dag(
    schedule="@daily",
    start_date=datetime(2024, 1, 1),
    catchup=False,
    default_args={
        "retries": 2,
        "retry_delay": timedelta(minutes=5),
        "execution_timeout": timedelta(minutes=30),
    },
)
def ingest():
    @task(retries=5, retry_exponential_backoff=True,
          max_retry_delay=timedelta(minutes=20))
    def call_rate_limited_api():
        ...

    @task(retries=0)  # expensive, not yet idempotent
    def rebuild_mart():
        ...

    call_rate_limited_api() >> rebuild_mart()

ingest()

go deeper

for a junior

Recall that retries is a count of extra attempts and retry_delay is a timedelta between them, and that both are usually set once in default_args for the whole DAG.

for a middle

Explain the arithmetic (retries + 1 executions), the up_for_retry state the instance sits in while waiting, and what retry_exponential_backoff plus max_retry_delay change about the spacing.

for a senior

Show judgment about which failures are worth retrying at all, and be ready to explain how retry counts delay your failure alert and why a non-idempotent task must be fixed before retries are safe.

for a principal

Own the org-wide default: what retry policy ships in the shared default_args, how it interacts with execution_timeout and alert latency, and how you stop teams from using retries to paper over unreliable tasks.

## What the parameters mean Every Airflow operator inherits a set of retry-related arguments from `BaseOperator`. The two you set most often are: - **`retries`** — an integer count of attempts *after* the first one. `retries=0` means the task gets exactly one shot. `retries=3` means up to four executions of that task instance. - **`retry_delay`** — a `datetime.timedelta` saying how long to wait between a failure and the next attempt. Both are ordinary keyword arguments on any operator, and both are commonly hoisted into `default_args` so every task in the DAG inherits them. ```python from datetime import timedelta from airflow.decorators import dag, task default_args = { "retries": 3, "retry_delay": timedelta(minutes=5), } ``` ## What actually happens on a failure When a task instance raises, the worker records the failure and the scheduler decides what state to put it in. If the instance's try number is still below `retries + 1`, it becomes **`up_for_retry`** and gets a `next_retry_datetime` computed from `retry_delay`. The DAG run remains `running` — a task in `up_for_retry` does not fail the run. When the delay elapses, the scheduler queues the instance again and the try number increments. Only after the final attempt fails does the instance become **`failed`**, and only then do downstream tasks react (typically going to `upstream_failed` under the default `all_success` trigger rule). This is why a retried task's log page shows several numbered attempts: each try writes its own log file, and the UI lets you page between them. Reading only the last try is a common way to miss the actual cause, because an early attempt may show a different error than the final one. ## Backoff A fixed `retry_delay` retries at a constant cadence. That is fine for a flaky network call, but it hammers a downstream service that is genuinely down. Airflow offers `retry_exponential_backoff=True`, which grows the wait between successive attempts rather than keeping it flat, and `max_retry_delay` (also a `timedelta`) to put a ceiling on how long the gap can get. The combination gives you fast recovery from a one-off blip without a task pinned in `up_for_retry` for hours. ```python @task( retries=5, retry_delay=timedelta(seconds=30), retry_exponential_backoff=True, max_retry_delay=timedelta(minutes=20), ) def call_flaky_api(): ... ``` ## What retries do *not* fix Retries only help when the failure is transient — a dropped connection, a throttled API, a node evicted mid-run. They do nothing for a deterministic failure such as a syntax error, a schema mismatch, or a missing permission; those simply burn every attempt and delay the alert by `retries × retry_delay`. That delay matters: a task with five retries five minutes apart takes almost half an hour to reach `failed` and fire whatever alerting hangs off failure. Sizing retries is therefore also a decision about how quickly you want to be told. The second thing retries do not fix is non-idempotent work. Airflow re-executes the whole task body on each attempt; it does not resume from where the code stopped. If your task appended rows, sent emails, or incremented a counter before it failed, the retry does it again. Making the task safe to re-run — deleting-then-inserting for the run's data interval, using a `MERGE`, writing to a run-scoped path, or keying an external call on the interval — is a precondition for turning retries on at all, not an optimisation. ## Related knobs worth knowing - **`execution_timeout`** — a `timedelta` after which the task instance is killed and treated as failed, which then feeds the retry machinery. Without it, a hung task can occupy a worker slot indefinitely and no retry ever happens, because the attempt never ends. - **`on_retry_callback`** — a callable fired on each retry, useful for counting retries as a metric rather than only alerting on the terminal failure. - **`depends_on_past`** and `wait_for_downstream` interact with failures at the DAG-run level: a permanently failed instance can block later runs of the same task. ## Where to set them Put a conservative default in `default_args` — a couple of retries with a few minutes' delay covers most infrastructure flakiness — and override per task where the failure profile differs. A task that calls a rate-limited third-party API wants more attempts with exponential backoff; a task that runs an expensive warehouse query you do not want executed twice wants very few, or zero until the query itself is made idempotent.

  • If a task has retries=3, how many times does its code actually execute in the worst case?
    Four. `retries` counts attempts *after* the initial run, so the total is `retries + 1`. This trips people up when they set `retries=1` expecting a single execution and get two. Each attempt writes its own numbered log in the task instance's log view.
  • Why can generous retries make an incident worse rather than better?
    Two reasons. Deterministic failures burn every attempt, so `retries=5` with a 10-minute `retry_delay` delays the failure alert by nearly an hour. And if the task is not idempotent, each attempt repeats whatever partial side effect the previous one left behind — duplicated rows, re-sent notifications, double-counted increments.
  • What stops a hung task from ever reaching the retry logic?
    Nothing, unless you set `execution_timeout`. Retries only fire when an attempt *ends* in failure; a task blocked on a socket read never ends, so it holds its worker slot indefinitely and no retry is scheduled. `execution_timeout` bounds the attempt, kills it, and lets the retry machinery take over.

saying these in an interview costs you the question

  • Says retries=1 means the task runs once total
  • Believes a retry resumes from where the code stopped
  • Turns on retries without checking the task is idempotent
  • Thinks a task in up_for_retry has already failed the DAG run
  • Sets huge retry counts on deterministic failures like schema errors

context

open as a page

How do you get alerted when an Airflow task fails, and what do the callback arguments give you?

level: middleimportance: must knowfreq 72%

basics

~20 s

Airflow 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.

open as a page

Which Airflow task instance states do you check first when a DAG run appears stuck?

level: seniorimportance: must knowfreq 66%

basics

~20 s

Read the task instance state in Airflow's grid view: scheduled or queued means no worker slot, running means the work is live, up_for_retry means it is waiting between attempts, and none or upstream_failed means it was never eligible. The state names the layer to investigate.

open as a page

In Airflow, what does setting an sla on a task do, and how does it differ from execution_timeout?

level: middleimportance: should knowfreq 54%

basics

~20 s

In Airflow an sla is a timedelta measured from the DAG run's data interval end; if the task has not succeeded by then, Airflow records an SLA miss and fires sla_miss_callback. It never interrupts the task, unlike execution_timeout, which kills a running attempt.

open as a page

What is a zombie task in Airflow, and what causes the scheduler to mark one?

level: seniorimportance: should knowfreq 58%

basics

~20 s

An Airflow zombie task is a task instance the metadata database still shows as running while no live worker process is heartbeating for it. The scheduler detects the missing heartbeats and fails the instance so its retries or downstream handling can proceed.

open as a page