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