🎉 75% of content is free forever — Unlock Premium from $10/mo →
CW
đŸ’ŧ Servicesâ„šī¸ Aboutâœ‰ī¸ ContactView Pricing Plansfrom $10

Error Handling and Retry Strategies in Apache Airflow

đŸŸĸ Free Lesson

Advertisement

Error Handling and Retry Strategies

Error Handling ArchitectureTask ExecutionRunning stateRetry LogicAutomatic retryTimeoutExecution limitsCallbackson{'_'}failure callbackSLA MonitoringSLA miss alertsRetry Delay Formuladelay = base{''}delay * exponential{''}backoff^retry{'_'}numberState Transitionsrunning {'->'} failed {'->'} up{''}for{''}retry {'->'} queuedAlert ChannelsEmail, Webhook, PagerDuty, SlackKey: retries=3, retry{'_'}delay=5min, timeout=3600s (typical config)

Architecture Diagram

Formal Definitions

Detailed Explanation

Retry Configuration

Airflow provides built-in retry mechanisms for handling transient failures.


Retry Parameters:

ParameterDescriptionExample
retriesNumber of retry attempts3
retry_delayTime between retriestimedelta(minutes=5)
retry_exponential_backoffDouble delay each retryTrue
max_retry_delayMaximum retry delay captimedelta(minutes=30)
execution_timeoutMax task execution timetimedelta(hours=1)
@task(
    retries=3,
    retry_delay=timedelta(minutes=5),
    retry_exponential_backoff=True,
    max_retry_delay=timedelta(minutes=30),
    execution_timeout=timedelta(hours=1),
)
def flaky_api_call():
    # Retries: 5min → 10min → 20min (capped at 30min)
    pass

Callback Functions

Callbacks are invoked on specific task state transitions.


Callback Types:

CallbackTriggerUse Case
on_failure_callbackTask failsSend alerts, log errors
on_success_callbackTask succeedsNotify, update external systems
on_retry_callbackTask retriesLog retry attempts
on_execute_callbackTask startsInitialize resources
def failure_callback(context):
    task_instance = context['task_instance']
    exception = context['exception']
    send_email(
        to=['team@example.com'],
        subject=f"Task Failed: {task_instance.task_id}",
        html_content=f"Exception: {exception}",
    )

SLA Configuration

SLAs define expected completion times for tasks.


SLA Parameters:

LevelSettingEffect
DAG-levelsla_miss_callbackInvoked when any task misses SLA
Task-levelsla=timedelta(hours=2)Per-task SLA

SLA Violation Condition: T_completion > T_sla + T_execution_date

@dag(
    sla_miss_callback=sla_miss_callback,
    default_args={'sla': timedelta(hours=2)},
)
def sla_example_dag():
    @task(sla=timedelta(hours=1))
    def critical_task():
        pass  # Must complete within 1 hour

Need Expert Airflow Help?

Get personalized tutoring, project support, or professional consulting.

Advertisement