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

Airflow 進階:Datasets、deferrable operators 與 executor 選型

· tech

#airflow#data-engineering

📑 目錄

這篇講三個進階功能。它們看似無關,其實有個共同點:都在解 Airflow 某個「死板」或「浪費」——排程只認時間、sensor 佔著 slot 空等、worker 一刀切。三個都是選配,但每個都能讓你的平台更聰明、更省。

從時間驅動到資料驅動:Datasets

跨 DAG 協調,最土的做法是用 cron 對時間:DAG A 排 02:00,DAG B 排 02:30——賭 A 半小時內跑完。但 A 一遲到,B 在 02:30 照樣跑,吃到舊資料甚至空的。Datasets(資料感知排程)把這件事翻轉成資料驅動:

時間驅動:cron 對時(脆弱) DAG A @ 02:00 賭 A 跑完了 DAG B @ 02:30(猜的) A 遲到 → B 照跑→ 吃到舊 / 空資料 ✗ 資料驅動:Dataset(精準) DAG A → 產出 sales Dataset: sales DAG B schedule=[sales] A 一更新 → B 立刻觸發,不早不晚 ✓
時間驅動是「賭上游好了」——DAG B 排在 A 後面半小時,A 一遲到就吃到舊資料。資料驅動是「上游告訴我好了」:DAG A 宣告它產出(produces)一個 Dataset,DAG B 的 schedule 直接設成那個 dataset——A 一更新它,B 立刻被觸發,不早不晚。這比舊的 cron 對時或 ExternalTaskSensor 精準太多,是從輪詢時間被事件驅動的轉變

落成程式碼,就是「A 掛 outlet、B 拿 dataset 當 schedule」:

from airflow.datasets import Dataset

sales = Dataset("s3://warehouse/sales")

@task(outlets=[sales])          # A:宣告「我這個 task 會更新 sales」
def build_sales(): ...

@dag(schedule=[sales], ...)     # B:sales 一被更新就自動觸發(不用寫 cron)
def downstream(): ...

別讓 sensor 佔著 slot 空等:deferrable operators

Sensor 那篇講過 sensor 用來「等外部條件」(等一個檔案、等一個外部 job)。但傳統 sensor 有個浪費:它整段等待時間都佔著一個 worker slot 在那反覆 poke。等的東西越多,slot 被空等吃得越乾淨:

傳統 sensor:佔 slot 空等 worker 池 slot 1:sensor 空等… slot 2:sensor 空等… slot 3:sensor 空等… N 個 sensor = N 個 slot 卡死 worker 被「等待」吃光,真 task 排不進 deferrable:丟給 triggerer worker 池(slot 釋放) 空出來 → 去跑真正的 task ✓ Triggerer(1 個 async 進程)用 asyncio 同時等上千個 等待幾乎免費;事件到了才喚回 worker
傳統 sensor 在等的整段時間,都佔著一個 worker slot 反覆 poke——10 個 sensor 就卡死 10 個 slot,真正要幹活的 task 反而排不進。Deferrable operator 把「等待」這件事 defer 出去,交給 triggerer(一個 async 進程,用 asyncio 同時等上千個),worker slot 立刻釋放去跑真 task,等到事件成立才把 task 喚回。等待從「昂貴地佔資源」變成「幾乎免費」

要用它,除了把 operator 換成 deferrable 版本(很多內建 sensor 有 deferrable=True 開關),還得確定叢集有跑 triggerer 這個進程(infra-airflow 提過它是常駐元件之一)。當你有一堆在等外部系統的 task,這個改動能把 worker 的有效產能翻好幾倍。

選對執行器:executor

最後一個「別一刀切」的地方是 executor——它決定 task 實際怎麼被執行。infra 那篇從部署角度深談過,這裡給一張「該選哪個」的速查:

Executortask 怎麼跑適合
Local在 scheduler 同機、以子行程跑開發、小專案、量不大
Celery丟進 broker(Redis/RabbitMQ),固定一群常駐 worker 搶著跑任務量穩定、想省開 pod 的延遲
Kubernetes每個 task 開一個 pod,跑完即刪任務起伏大、要徹底隔離與彈性擴縮

判準跟 Spark 的 executor 一樣落在 stateless 那條軸:任務量穩→Celery(固定池省延遲),任務量爆起爆落→Kubernetes(一 task 一 pod、算完即刪、不養閒置)。 別從頭到尾只用同一種。

反思

從「按時間猜」到「按資料反應」,是可靠性的一次躍遷

Datasets 讓我體會到一個通用的道理:time-driven 是猜,event-driven 是知道。 cron 排程的本質,永遠是「賭上游那時候好了」——賭錯就吃舊資料。而 Datasets 把它翻成「上游好了會告訴我」,那份不確定性就消失了。這個從輪詢時間被事件驅動的轉變,不只在 Airflow——webhook 之於輪詢 API、reactive 之於 polling、訊息驅動之於定時掃表,都是同一次躍遷。我現在看到任何「排個固定時間去猜對方好了沒」的設計,都會反射性地問一句:能不能改成『對方好了就通知我』? 那幾乎總是更準、也更省。

「等待」不該佔資源——這是個到處適用的效率洞見

deferrable operator 點破了一個我以前沒意識到的浪費:一個 sensor task,90% 的時間都在「等」,卻整段佔著一個 worker。把「等」和「做」分開、讓等待幾乎免費——這個念頭一旦有了,你會發現它到處都是:非阻塞 I/O、async/await、Redis 的 epoll、作業系統的中斷 vs 輪詢……全是同一個智慧:用一個執行緒空轉去等一件事,是最笨的做法;正確的做法是「掛起、讓出、事件來了再喚醒」。 Airflow 的 triggerer 就是把這個 OS 級的老智慧,搬到了工作流程這一層。看懂一個,你就看懂了一整類效率問題的解法。

收尾整個系列:成熟不是會用最多功能,是知道何時還不需要

這三個功能都很香,但我要強調的收尾,恰恰是它們都是選配。沒有跨 DAG 協調的痛,別急著上 Datasets;sensor 不多、slot 不緊,deferrable 是過度工程;單機跑得動,別為了潮而上 KubernetesExecutor。這呼應了整個 Airflow 系列、乃至我對工具一貫的態度——先確認痛點,再上重武器。回頭看這一路:從跑起第一個 DAG、搞懂排程區間、學會可靠性測試部署,到這篇的進階武器——真正的成熟,從來不是會用最多功能,而是清楚每個功能在解什麼痛、以及什麼時候你還不需要它。 一個 DAG 從能跑、到可靠、到聰明,靠的不是堆功能,是每一步都問對「我現在真正缺的是什麼」。Airflow 系列,就收在這句話上。