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——因為:
@task 裡,只在執行時跑。這也解釋了 infra 那篇說的「scheduler 猛查 DB」,你的 top-level 常是幫兇一句話記住:DAG 檔案描述「要做什麼」,不「現在就做」。 重活放進 task。
怎麼測 DAG:三個層次,由便宜到貴
DAG 也是程式碼,一樣要測——而且投報率最高的,是那個最無聊的:
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 都用 KubernetesPodOperator、SparkSubmitOperator 去跑,邏輯不就一定「包在 operator 裡」,還怎麼抽成純函數測?這其實是把兩種 operator 混為一談了:
- 執行型(
PythonOperator/@task):在 Airflow worker 行程裡真的跑你的 Python——邏輯就在 callable 裡,這才是要抽出來的那一塊。 - 提交型(
KubernetesPodOperator、SparkSubmitOperator):它們不執行你的邏輯,只負責「用哪個 image、帶什麼參數、丟去哪個叢集」然後 submit。你的邏輯根本不在 Airflow 裡,在 image / Spark job 的獨立 codebase。
所以「跑在 K8s / Spark」反而是邏輯與編排分離做得最徹底的形式——Airflow 只當交通指揮(編排歸編排、運算歸運算)。邏輯放哪裡,就在哪裡測:
| 邏輯放哪 | 誰執行 | 在哪測 |
|---|---|---|
@task / PythonOperator callable | Airflow 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-airflow | metadata 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 才算真的長大成人。