Airflow 可靠性實戰:冪等、重試、SLA 與告警

· tech

#airflow#data-engineering#reliability

📑 目錄

前面幾篇教你把 DAG 寫出來、排對區間。但 Production 的 DAG 是會在半夜出事的——來源系統延遲、網路抖一下、下游資料庫重啟。這篇講怎麼讓 DAG 扛得住失敗、能自己救、真救不了會大聲叫。我把可靠性拆成三層防線:冪等(重跑安全的地基)→ retries(自動自癒)→ SLA / 告警(兜底叫人)。三層缺一不可,而且順序不能顛倒——因為上面兩層,全都踩在「冪等」這塊地基上。

地基:冪等——沒有它,retry 只會放大災難

排程那篇說過:catchup、backfill、手動重跑,全都在「重跑同一段區間」。而 retries 也是重跑。這代表你的 task 注定會被跑不只一次,所以第一個問題從來不是「它會不會失敗」,而是「它失敗重跑,會不會出事」:

非冪等:INSERT 附加 第 1 次:INSERT 6/18 → 中途失敗 重試:又 INSERT 6/18 一次 6/18 資料重複兩份 ✗重試把一次故障放大成資料污染 冪等:覆寫分區 第 1 次:覆寫 6/18 分區 → 失敗 重試:再覆寫 6/18 分區一次 6/18 分區一致 ✓跑幾次結果都一樣,重試安全
同一個 task 被重試,冪等與否結果天差地遠。非冪等(對表 INSERT 附加)重試一次,資料就多一份——retry 反而把一次小故障放大成資料污染。冪等(覆寫「這段區間」的分區)不管跑幾次,結果都是同一份。冪等不是進階技巧,是 retries / backfill 能安全成立的前提

實作冪等的核心心法,就是排程那篇說的「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):

running up_for_retry等 retry_delay 失敗 重試(還有次數) 次數用完 failed 🔔 on_failure_callback→ Slack / PagerDuty 排程後超過 SLA 還沒完成 → SLA miss(獨立於失敗) 🔔 sla_miss_callback「太慢」也要叫 告警只在重試都用完後才響 → 暫時性故障被 retry 默默吸收,不吵人
一個 task 失敗後先進 up_for_retry、等 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 的半夜。