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

#airflow

時間驅動:cron 對時(脆弱) DAG A @ 02:00 賭 A 跑完了 DAG B @ 02:30(猜的) A 遲到 → B 照跑→ 吃到舊 / 空資料 ✗ 資料驅動:Dataset(精準) DAG A → 產出 sales Dataset: sales DAG B schedule=[sales] A 一更新 → B 立刻觸發,不早不晚 ✓

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

· tech · 約 5 分鐘 · 📚 Airflow 學習筆記 #9

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

#airflow#data-engineering

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 執行時跑一次

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

· tech · 約 7 分鐘 · 📚 Airflow 學習筆記 #8

DAG 寫好了、也可靠了,最後一哩是工程紀律:怎麼確定我的改動不會弄壞它,又怎麼把它安全送上 Production。 這篇講測試(把壞 DAG 擋在 merge 前)、部署(DAG 檔怎麼上環境、se…

#airflow#data-engineering#testing

非冪等:INSERT 附加 第 1 次:INSERT 6/18 → 中途失敗 重試:又 INSERT 6/18 一次 6/18 資料重複兩份 ✗重試把一次故障放大成資料污染 冪等:覆寫分區 第 1 次:覆寫 6/18 分區 → 失敗 重試:再覆寫 6/18 分區一次 6/18 分區一致 ✓跑幾次結果都一樣,重試安全

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

· tech · 約 5 分鐘 · 📚 Airflow 學習筆記 #7

前面幾篇教你把 DAG 寫出來、排對區間。但 Production 的 DAG 是會在半夜出事的——來源系統延遲、網路抖一下、下游資料庫重啟。這篇講怎麼讓 DAG 扛得住失敗、能自己救、真救不了會大聲…

#airflow#data-engineering#reliability

看似無狀態的組件,真狀態藏在一顆 metadata DB Scheduler解析 DAG、排 task無狀態・可重啟 Webserver那個 UI無狀態 Worker × N執行 task可多開・可重啟 Metadata DB(Postgres)真狀態・命門 組件掛了都能換;DB 掉了 = 整個系統的記憶歸零(哪些跑過、誰在跑、誰失敗) 連多個 scheduler 的 HA,都靠對這顆 DB 上鎖(row lock)來協調

Airflow:排程器、worker 與那個藏起來的狀態

· tech · 約 6 分鐘 · 📚 從 Infra 角度看資料工具 #7

上一篇埋了個伏筆:每個系統都有一塊逃不掉的狀態,認出它就掌握了命門。 Airflow 是這句話最好的示範。它表面上全是看似無狀態、可重啟的組件——scheduler、webserver、worker,…

#infrastructure#airflow

K8s Control Plane · Scheduler 依 resource request / affinity 決定 pod 落哪個 node Node 1 on-demand 池(穩定) Node 2 spot 池 Node 3 spot 池 Airflow Scheduler Airflow Webserver Airflow Triggerer Metadata DBPersistentVolume(有狀態) Spark Driver Spark Executor Spark Executor Spark Executor Spark Executor Driver 申請的 executor 由 Scheduler 散到各 node Executor 讀取來源 / 寫回結果 Business DB叢集外的營運資料源 / 輸出 Airflow pod Spark Driver Spark Executor Metadata DB Business DB

Airflow + Spark 跑在 K8s 上:不同的 node 怎麼跑不同的 pod

· tech · 約 11 分鐘 · 📚 Kubernetes 學習筆記 #8

Airflow 負責「什麼時候、用什麼順序」跑作業,Spark 負責「把大資料算完」。當這兩個東西都搬到 Kubernetes 上,最常見的困惑是:到底有哪些東西在跑、它們各自是不是一個 pod、又被…

#kubernetes#airflow#spark#data-engineering

branch full_reload ✓ incremental ⊘ skip join none_failed_…

Airflow 複雜流程控制:branching、trigger rules、TaskGroup、動態任務

· tech · 約 5 分鐘 · 📚 Airflow 學習筆記 #6

到目前為止,我們的 DAG 都是「固定的幾個 task、照線跑完」。但真實流程沒這麼乖:有時要依條件走不同分支、某個 task 失敗了清理工作還是得跑、幾十個相似 task 要收整齊、甚至執行期才知道…

#airflow#data-engineering#dag-design

DAG task @task Operator 做一件事 Hook 連線 client 外部系統 DB / S3 / API Connection conn_id:帳密 + 端點 提供帳密

Airflow 怎麼連外部系統:Provider、Operator、Hook、Sensor

· tech · 約 6 分鐘 · 📚 Airflow 學習筆記 #5

前面幾篇都在講 Airflow 的內部機制 —— 排程、任務間傳資料。但真正的 pipeline 一定得碰外面的世界:查一個資料庫、丟檔案到 S3、打一支 HTTP API、等一個檔案到齊。這篇講 A…

#airflow#data-engineering#integration

task A XCom 訊息板 (metadata DB) task B push pull 只放「小東西」:id、筆數、旗標、檔案路徑… 大資料請寫到 S3 / 檔案,XCom 只傳「路徑」

Airflow 任務間怎麼傳資料:XCom、TaskFlow 進階與 params

· tech · 約 4 分鐘 · 📚 Airflow 學習筆記 #4

上一篇講完排程,這篇處理一個很實際的問題:一個 task 算出來的東西,怎麼交給下一個 task? Airflow 的答案是 XCom,而 TaskFlow 把它包得幾乎看不見。順帶也講 params…

#airflow#data-engineering#xcom

6/18 00:00 6/19 00:00 6/20 00:00 資料區間 = 6/18 資料區間 = 6/19 6/18 的 run 在這觸發 6/19 的 run

Airflow 排程的真相:data interval、catchup 與 backfill

· tech · 約 5 分鐘 · 📚 Airflow 學習筆記 #3

上一篇我們把 DAG 跑起來了,還特別把 catchup 設成 False「保平安」。這篇就來拆解 Airflow 最反直覺、也最多人卡住的部分:排程到底怎麼運作。搞懂這一段,你才真的會用 Airfl…

#airflow#data-engineering#scheduling

./dags/*.py 你寫的 DAG Scheduler Web UI :8080 Worker 瀏覽器 掛載/解析 顯示 派送任務

跑起第一個 Airflow:Docker 環境 + 你的第一個 DAG

· tech · 約 4 分鐘 · 📚 Airflow 學習筆記 #2

上一篇把 Airflow 的概念與架構講清楚了;這篇動手把它在本機跑起來,寫出第一個 DAG,並在 Web UI 上看著它跑完。目標很單純:從零到「我親眼看到自己的 DAG 在介面上變綠」。 方式 指…

#airflow#data-engineering#docker