पाठ 19 / 25
Retries, Timeouts और Callbacks
Backoff के साथ retries, execution timeouts और failure callbacks configure करें।
विफलता की योजना
कई विफलताएँ अस्थायी होती हैं (network की रुकावट, rate limits), इसलिए retries और retry_delay तय करें, और प्रयासों के बीच ज़्यादा इंतज़ार के लिए retry_exponential_backoff=True। execution_timeout रखें ताकि अटका task हमेशा चलने की जगह मार दिया जाए, और DAG-स्तर का dagrun_timeout। Alert भेजने (Slack, email, pager) के लिए DAG, task और logs के link के साथ on_failure_callback उपयोग करें। Retry सिर्फ़ idempotent tasks के लिए सुरक्षित है, इसीलिए डिज़ाइन मायने रखता है। साझा settings default_args में रखें और प्रति task बदलें।
ज़ोर से विफल हों, चुपचाप सुधरें
Retries, timeouts, alerts और क्षमता नियंत्रण pipelines को भरोसेमंद रखते हैं।
Default args
उदाहरण मान। Callback function को task instance और exception वाली context dictionary मिलती है।
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,
}Backoff अनुसूची, चलाकर
मैंने तीन retries के लिए घातीय प्रतीक्षा (60 s, दोगुनी, 900 s पर सीमित) का सरल मॉडल चलाया। Airflow का अपना सूत्र jitter और ब्योरे जोड़ता है; यह आकार दिखाता है।
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]
त्वरित जाँच: Task को retry करना कब सुरक्षित है?
- सिर्फ़ बिना तारीख़ वाले tasks के लिए
- कभी नहीं
- सिर्फ़ जब वह बिना keys के rows append करता हो
- जब task idempotent हो
Answer
जब task idempotent हो — Idempotent task दोहराने से डेटा दोहरा या ख़राब नहीं हो सकता।