Lesson 19 / 25

Retries, Timeouts and Callbacks

Configure retries with backoff, execution timeouts and failure callbacks.

Plan for failure

Many failures are temporary (network blips, rate limits), so set retries and retry_delay, and retry_exponential_backoff=True to wait longer between attempts. Set execution_timeout so a hung task is killed instead of running forever, and a DAG-level dagrun_timeout. Use on_failure_callback to send an alert (Slack, email, a pager) with the DAG, task and a link to the logs. Retrying is safe only for idempotent tasks, which is why design matters. Put shared settings in default_args and override per task.

Fail loudly, recover quietly

Retries, timeouts, alerts and capacity controls keep pipelines dependable.

Three controls: retry, limit, alert.
Figure 6.1 — Retry, limit and alert.

Default args

Illustrative values. The callback function receives a context dictionary with the task instance and exception.

from datetime import timedelta

def alert_on_failure(context):
    ti = context["task_instance"]
    send_slack(f"FAILED: {ti.dag_id}.{ti.task_id} run={context['run_id']} log={ti.log_url}")

default_args = {
    "retries": 3,
    "retry_delay": timedelta(minutes=2),
    "retry_exponential_backoff": True,
    "max_retry_delay": timedelta(minutes=15),
    "execution_timeout": timedelta(hours=1),
    "on_failure_callback": alert_on_failure,
}

A backoff schedule, run

I ran a simple model of exponential waits (60 s, doubling, capped at 900 s) for three retries. Airflow's own formula adds jitter and details; this shows the shape.

def retry_schedule(delay_s=60, retries=3, factor=2, cap=900):
    return [min(cap, delay_s * factor**i) for i in range(retries)]

print(retry_schedule())

Output:

[60, 120, 240]

Quick check: When is retrying a task safe?

  • Only for tasks without dates
  • Never
  • Only when it appends rows without keys
  • When the task is idempotent
Answer

When the task is idempotent — Repeating an idempotent task cannot duplicate or corrupt data.