🌐 This page hasn't been translated yet — showing the original Chinese. Translated posts

Airflow 測試與部署:別讓一個 typo 弄垮整包 DAG

· tech

#airflow#data-engineering#testing

📑 目錄

DAG 寫好了、也可靠了,最後一哩是工程紀律:怎麼確定我的改動不會弄壞它,又怎麼把它安全送上 Production。 這篇講測試(把壞 DAG 擋在 merge 前)、部署(DAG 檔怎麼上環境、self-host vs managed),以及一個最要命、幾乎每個新手都踩過的 Production 陷阱——DAG 檔案不是「執行一次」,是「一直被解析」的。

最要命的陷阱:DAG 檔案一直被 scheduler 解析

新手常把 DAG 檔當普通腳本寫,在檔案最上層(top-level)直接連資料庫、打 API、跑一段重運算。這在本機跑一次沒事,上了 Production 卻會悄悄拖垮整個 scheduler——因為:

Scheduler每隔幾秒 parse 一次所有 DAG 檔 parseparse ✗ top-level 放重活 df = query_db() # 模組層 cfg = requests.get(url) 每次 parse 都跑這些 → scheduler 被反覆拖垮、DAG 變慢 ✓ 檔案輕,重活進 task @task def load(): query_db() top-level 只有 DAG 結構 parse 秒過 重活只在 task 執行時跑一次
關鍵事實:scheduler 每隔幾秒就重新解析(import)一次所有 DAG 檔,好知道有沒有新 DAG、排程有沒有變。所以檔案最上層的任何程式碼,都會被反覆執行在 top-level 連 DB、打 API,等於叫 scheduler 每幾秒幫你 query 一次,整個排程被你自己拖垮。正解是讓 DAG 檔很輕——top-level 只放結構定義,真正的重活一律放進 @task 裡,只在執行時跑。這也解釋了 infra 那篇說的「scheduler 猛查 DB」,你的 top-level 常是幫兇

一句話記住:DAG 檔案描述「要做什麼」,不「現在就做」。 重活放進 task。

怎麼測 DAG:三個層次,由便宜到貴

DAG 也是程式碼,一樣要測——而且投報率最高的,是那個最無聊的:

① import 驗證(最便宜、擋最多)所有 DAG import 無錯、無 cycle ② task 邏輯單元測試邏輯抽成純函數,直接測 ③ dag.test() 本地整跑不碰 metadata DB CI:每個 PR 都跑這三層+ lint,綠燈才能 merge merge → deploy DAG 到環境git-sync / S3 / 烤進 image 越下面越貴、越少 越上面越便宜、越該先做 光是「所有 DAG 都 import 得起來」就擋掉最多事故
測 DAG 分三層,投報率由高到低:①import 驗證最無聊也最值——一個 pytest 把所有 DAG import 一遍、斷言無語法錯無循環,就能擋掉 Production 最常見的事故(一個 typo 讓整包 DAG 從 UI 上消失);②單元測試要先把商業邏輯抽成純函數(別埋在 operator 裡)才好測;dag.test() 本地把整個 DAG 跑一遍、不碰 metadata DB。這三層擺進 CI,綠燈才 merge,壞 DAG 就進不了 Production

那個最該先寫的 import 測試,其實非常短:

# 擋掉最多事故的一個測試:所有 DAG 都 import 得起來、沒循環
from airflow.models import DagBag

def test_dags_load_without_errors():
    dag_bag = DagBag(include_examples=False)
    assert dag_bag.import_errors == {}, f"DAG 匯入失敗:{dag_bag.import_errors}"

而要單元測 task 邏輯,關鍵是把邏輯抽成純函數——別把運算埋在 operator 裡,那樣得起一整個 Airflow 才測得到:

# 邏輯抽出來,pytest 直接測,不需要 Airflow
def compute_summary(rows: list[dict]) -> dict: ...

@task
def summarize(**ctx):
    return compute_summary(fetch(ctx["ds"]))   # task 只是薄薄一層黏合

第三層,dag.test()(Airflow 2.5+)能在本機把整個 DAG 跑一遍、不需要 scheduler、不寫 metadata DB,拿來在 CI 或本地做端到端驗證很順手。

但跑在 K8s / Spark 上呢?邏輯根本不在 Airflow 裡

有人會問:實務上 task 都用 KubernetesPodOperatorSparkSubmitOperator 去跑,邏輯不就一定「包在 operator 裡」,還怎麼抽成純函數測?這其實是把兩種 operator 混為一談了:

  • 執行型(PythonOperator / @task):在 Airflow worker 行程裡真的跑你的 Python——邏輯就在 callable 裡,這才是要抽出來的那一塊
  • 提交型(KubernetesPodOperatorSparkSubmitOperator):它們不執行你的邏輯,只負責「用哪個 image、帶什麼參數、丟去哪個叢集」然後 submit。你的邏輯根本不在 Airflow 裡,在 image / Spark job 的獨立 codebase。

所以「跑在 K8s / Spark」反而是邏輯與編排分離做得最徹底的形式——Airflow 只當交通指揮(編排歸編排、運算歸運算)。邏輯放哪裡,就在哪裡測:

邏輯放哪誰執行在哪測
@task / PythonOperator callableAirflow worker抽成純函數,pytest 直接測
Spark job(SparkSubmitOperator 提交)Spark 叢集在 Spark job 自己的 repo,用 pyspark local session 測
容器(KubernetesPodOperator)K8s pod在容器 app 自己的 repo 測
DAG 的接線本身——DagBag import + 斷言 operator 的 image / args / 依賴設對

一句話:「別包邏輯進 operator」不是「不准用 KubernetesPodOperator」——前者是叫你別把 transform 寫死在 PythonOperator 的 callable;後者本來就是把邏輯推去容器 / Spark,方向完全一致。用提交型 operator 時,DAG 這層只要測「有沒有接對」,不必真的把 Spark 跑起來。

部署:DAG 檔怎麼上環境,以及 self-host vs managed

Airflow 的「部署」,本質就是把 DAG 檔案送到執行環境,常見三種:git-sync(一個 sidecar 定期把 git repo 拉到 DAG 目錄)、丟物件儲存(MWAA 就是把 DAG 放 S3)、烤進映像檔(Astronomer / K8s 上,把 DAG 打包成 image 一起部署,最可控)。至於整套 Airflow 要不要自己養,是另一個更大的決策:

Self-host(Helm on K8s)Managed(MWAA / Astronomer)
你要顧scheduler、[[infra-airflowmetadata DB]]、升級、擴縮
成本授權免費,但吃維運人力付費,換掉那座維運的山
適合有 K8s 底子、要省錢或要深度客製想把心力放在資料、不想養基礎設施

這個選擇跟我對所有基礎設施的態度一致:別只看帳面授權費,要把「維運人力」算進去。Airflow 背後那顆 metadata DB、那個要顧的 scheduler、那些升級,都是真實的人力成本——managed 是付錢把這座山租給別人。團隊小、又沒人想當 Airflow 管理員時,managed 幾乎總是划算的。

反思

「DAG 檔是被反覆解析的」,是我學 Airflow 最貴的一課

這件事沒人一開始會告訴你,你得親手把 scheduler 拖垮一次才會刻進骨子裡。想通「這個檔案每幾秒就被 import 一次」之後,我寫 DAG 的方式徹底變了:檔案最上層只留結構、像一張宣告,所有會連外、會算、會慢的東西一律鎖進 task。這其實是一個更通用的直覺——先搞清楚你寫的程式碼「多久被執行一次」,再決定把什麼放哪裡。 放錯層級的重活,是很多系統莫名其妙變慢的根源,而它在 Airflow 上特別致命,因為那個「反覆執行」是隱形的,藏在 scheduler 的解析迴圈裡。

可靠性常常不是靠聰明,是靠把最笨的檢查自動化

測 DAG 這件事,最反直覺的收穫是:投報率最高的測試,是那個最無聊的 import 驗證。 它不驗證任何業務邏輯,只確認「所有 DAG 都還 import 得起來、沒有 typo、沒有循環」——但就這個最低標,擋掉了我看過最多的 Production 事故:某人改一行、一個 import 打錯,結果整包 DAG 從 UI 上消失,排程默默停擺,沒人發現。這讓我對「可靠性」的理解更務實了:它往往不是來自多精巧的測試策略,而是來自把那些「顯而易見卻沒人做」的笨檢查,變成 CI 裡自動跑的一關。聰明的測試是加分,笨的檢查自動化是及格線——而多數團隊連及格線都沒守住。

部署要算的是「維運人力」,不是授權費

self-host vs managed 的抉擇,把我對成本的看法擦亮了。Airflow 開源免費,但「免費」只是授權欄那一格——它背後的 metadata DB 要有人備份、scheduler 要有人盯、版本要有人升,這些都是看不見卻真實的人力帳單。我看過團隊為了省 managed 的月費,自建 Airflow,結果一個工程師半條命都在修它,那省下的錢遠不夠付這個隱形成本。所以我現在評估任何「要不要自己架」時,第一個問的不是「授權多少錢」,而是「這座山,我們有沒有人、願不願意養」。這跟 先確認痛點再上重武器是同一種務實:工具的真實成本,從來不在它的價目表上。這也是 Airflow 主線的收尾——會寫、可靠、測得住、部署得起,一個 DAG 才算真的長大成人。