पाठ 16 / 25
Idempotent, Partitioned Tasks
हर run का output निश्चित, तारीख़-आधारित स्थान पर लिखें ताकि re-runs जोड़ने की जगह overwrite करें।
वही input, वही output, कितनी भी बार
Tasks विफल होते हैं, retry होते हैं और पिछली तारीख़ों के लिए दोबारा चलाए जाते हैं, इसलिए हर task idempotent होना चाहिए: एक ही logical date के लिए उसे दो बार चलाना वही अंतिम स्थिति छोड़े जो एक बार चलाना। तकनीकें: run की तारीख़ से निकले path या partition (.../dt=2026-10-01/) में लिखें और उसे overwrite करें; अंधे INSERT की जगह तारीख़ की rows के लिए MERGE/upsert या DELETE फिर INSERT उपयोग करें; अस्थायी स्थान पर लिखकर atomically swap करें; de-duplication key के बिना कभी append न करें। Idempotency ही retries और backfills को सुरक्षित बनाती है।
दो बार चलाना सुरक्षित
अच्छी pipelines किसी भी तारीख़ के लिए बिना डेटा दोहराए या खोए दोबारा चलाई जा सकती हैं।
निश्चित output path, चलाकर
मैंने यह चलाया: path सिर्फ़ logical date पर निर्भर है, इसलिए retry या backfill ठीक उसी जगह लिखता है और पुरानी फ़ाइल बदल देता है।
def out_path(ds):
return f"s3://lake/orders/dt={ds}/part.parquet"
print(out_path("2026-10-01"), out_path("2026-10-01") == out_path("2026-10-01"))
Output:
s3://lake/orders/dt=2026-10-01/part.parquet True
तारीख़ के लिए delete-फिर-insert (SQL)
एक ही तारीख़ के लिए इसे दो बार चलाने पर उस दिन की rows की एक ही copy रहती है। {{ ds }} Airflow द्वारा templated है।
BEGIN;
DELETE FROM analytics.daily_orders WHERE order_date = '{{ ds }}';
INSERT INTO analytics.daily_orders
SELECT order_date, customer_id, SUM(amount) AS revenue
FROM staging.orders
WHERE order_date = '{{ ds }}'
GROUP BY order_date, customer_id;
COMMIT;त्वरित जाँच: Task idempotent क्यों होना चाहिए?
- क्योंकि Airflow files लिखना मना करता है
- उसे धीमा करने के लिए
- ताकि retries और re-runs डेटा दोहराएँ या बिगाड़ें नहीं
- तारीख़ों के उपयोग से बचने के लिए
Answer
ताकि retries और re-runs डेटा दोहराएँ या बिगाड़ें नहीं — सुरक्षित दोहराव ही orchestrators को भरोसे से retry और backfill करने देता है।