Airflow 可靠性實戰:冪等、重試、SLA 與告警
· tech
#airflow#data-engineering#reliability
📑 目錄
前面幾篇教你把 DAG 寫出來、排對區間。但 Production 的 DAG 是會在半夜出事的——來源系統延遲、網路抖一下、下游資料庫重啟。這篇講怎麼讓 DAG 扛得住失敗、能自己救、真救不了會大聲叫。我把可靠性拆成三層防線:冪等(重跑安全的地基)→ retries(自動自癒)→ SLA / 告警(兜底叫人)。三層缺一不可,而且順序不能顛倒——因為上面兩層,全都踩在「冪等」這塊地基上。
地基:冪等——沒有它,retry 只會放大災難
排程那篇說過:catchup、backfill、手動重跑,全都在「重跑同一段區間」。而 retries 也是重跑。這代表你的 task 注定會被跑不只一次,所以第一個問題從來不是「它會不會失敗」,而是「它失敗重跑,會不會出事」:
實作冪等的核心心法,就是排程那篇說的「run 只吃自己那段區間、並覆寫它」,而不是無腦附加:
@task
def load(**context):
ds = context["data_interval_start"].strftime("%Y-%m-%d")
# ✓ 冪等:先清掉這段區間、再寫入(或直接覆寫分區 / MERGE upsert)
conn.execute(f"DELETE FROM sales WHERE dt = '{ds}'")
conn.execute(f"INSERT INTO sales SELECT * FROM staging WHERE dt = '{ds}'")
# ✗ 非冪等寫法:INSERT INTO sales ...(沒先清,重跑就疊加)
寫物件儲存也一樣:寫到 s3://.../dt={{ ds }}/ 這種由區間決定的確定性路徑、整個覆蓋,而不是每次 append 一個新檔。把「重跑」設計成安全的,後面的 retries 才敢放手讓它自動跑。
第一層自癒:retries
地基穩了,就能讓 Airflow 自動處理暫時性故障——來源慢一秒、連線被重置這種,重試一下多半就過了。設在 default_args 讓整個 DAG 的 task 共用:
from datetime import timedelta
default_args = {
"retries": 3, # 失敗自動重試 3 次
"retry_delay": timedelta(minutes=5), # 每次間隔 5 分鐘
"retry_exponential_backoff": True, # 改成指數退避:5、10、20 分…
"max_retry_delay": timedelta(hours=1), # 退避上限
"execution_timeout": timedelta(minutes=30), # 跑超過 30 分就砍掉(防卡死)
}
兩個重點:一是 execution_timeout——一個 task 卡死(連線 hang、等一個永遠不來的檔案)比它失敗更可怕,因為它會一直佔著 slot、無聲無息,設個上限讓它「該死就死、然後進重試」。二是要認清 retries 只救得了暫時性故障;如果是程式邏輯錯、SQL 打錯,重試三次只是把同一個錯誤演三遍,白費五分鐘還延誤告警。分清「暫時 vs 永久」,別指望 retry 修永久的錯。
慢,也是一種故障:SLA
有個常被忽略的狀況:task 沒失敗,但太慢。一個每天早上八點該好的報表,今天拖到中午才出——對下游來說,跟壞了沒兩樣。Airflow 用 SLA 捕捉這種「準時性」故障:
@task(sla=timedelta(hours=2)) # 排程後 2 小時內沒跑完 = SLA miss
def daily_report():
...
SLA miss 是獨立於「失敗」的一條線:task 最後可能還是成功了,但它遲到了,Airflow 會記一筆 SLA miss、觸發 sla_miss_callback。把「太慢」也當成要被告警的可靠性事件,是成熟資料平台跟堪用平台的分界——可用性從來不只是「有沒有成功」,還有「有沒有及時」。
真的救不了:告警要大聲
當 retries 用完、task 真的 failed,或發生 SLA miss,這一刻必須有人知道。Airflow 用 callback 把這些時刻接出去(通常接到 Slack / PagerDuty):
retry_delay 再重跑,重複到次數用完才真正 failed,這時才觸發 on_failure_callback 發告警。這個設計的精妙在:暫時性故障被 retry 默默吸收,人只會被「真的救不了」的事叫醒。SLA miss 則是另一條獨立的線——沒失敗、但遲到了,一樣要叫。太吵的告警等於沒告警,好的可靠性讓警報「該響才響」callback 接出去長這樣:
def alert_to_slack(context):
ti = context["task_instance"]
send_slack(f"🔥 {ti.dag_id}.{ti.task_id} 失敗於 {context['ds']}")
default_args = {
"retries": 3,
"on_failure_callback": alert_to_slack, # 重試全用完、真 failed 才呼叫
# "on_retry_callback": ... # 想連每次重試都知會也可以(通常不必)
}
反思
冪等是可靠性的地基,不是進階選項
我對 Airflow 可靠性最深的體會,是它的一切都建立在「重跑安全」上。retries、catchup、backfill、災難重跑——這些便利功能,全都預設你的 task 冪等;一旦不冪等,retries 反而是最危險的東西,它會忠實地把一次小故障,自動放大成三份重複的髒資料。所以我寫每個 task 的第一個念頭,已經不是「怎麼讓它成功」,而是「它被跑第二次、第三次,會不會壞事」。這跟 把 now() 趕出 task、跟 SRE 那篇講的 exactly-once 是同一條紀律的不同臉:在一個什麼都會重跑的系統裡,冪等不是加分項,是入場券。
好的自癒,是把暫時性故障變成「非事件」
retries 教我一件關於告警的事:不是每個失敗都該叫醒人。網路抖一下、來源晚幾秒,這種暫時性故障,系統自己 retry 兩次就過去了,根本不該驚動任何人。真正該把人從床上挖起來的,是「retry 都用完了還救不回來」的事。所以我很在意告警的時機——讓它只在重試耗盡後才響,而不是每次小失敗都轟炸。這背後是一個更大的原則:告警的價值跟它的準確度成正比,跟它的數量成反比。 一個每天誤報十次的頻道,真出事時沒人會看;好的可靠性設計,是讓警報「該響才響」,這樣它一響,大家才會當真。這也是 SRE 監控那套的核心。
「慢」也是故障——可用性不只是「有沒有成功」
SLA 這個機制,把我對「可靠」的定義撐大了。以前我只盯著「有沒有跑成功」,直到理解:一個每天準時的報表,今天拖到中午才好,對等著它做決策的下游,傷害跟直接失敗沒兩樣。及時,本身就是一種正確。 從此我看資料管線的健康,不只看成功率,也看「準時率」——有沒有在下游需要它之前完成。把時效也納進監控與告警,是我認為「堪用的資料平台」跟「可信任的資料平台」之間,那條最關鍵的線。可靠性 = 冪等的地基 + retries 的自癒 + SLA 與告警的兜底,三層一起,你的 DAG 才真的扛得住 Production 的半夜。