TaskFlow API

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

傳統作法

def extract_data():
    return {"users": 150, "events": 4200}

with DAG("etl_pipeline") as dag: t1 = PythonOperator(task_id="extract", python_callable=extract_data) t2 = PythonOperator(task_id="summary", python_callable=print_summary)
t1 >> t2
使用 Airflow 建置資料管線

TaskFlow 作法

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
使用 Airflow 建置資料管線

為什麼用 TaskFlow?

 

  • 更少樣板碼,以裝飾器取代 operator 實例
  • 隱含相依,回傳值會自動串接任務
  • 可讀性高,Dag 就像一支 Python 指令稿
  • 傳統 operators 仍可用於 provider 整合

TaskFlow API - visual

使用 Airflow 建置資料管線

Airflow 介面:Grid 檢視

Airflow Grid view

 

  • 每一欄是一次 Dag 執行
  • 每一列是一個 任務
  • 顏色表示狀態:綠色=成功,紅色=失敗
  • 點選方格可查看 記錄與細節
使用 Airflow 建置資料管線

Airflow 介面:Graph 檢視

Airflow Grid view

 

  • 聚焦於單次 Dag 執行
  • 顯示 相依關係與平行執行
  • 可在設定中調整細節
使用 Airflow 建置資料管線

Airflow 介面:Dag 版本控管

Dag version indicator in Airflow UI

 

  • 自動追蹤 結構變更
  • 每次執行會連到當時 啟用的版本
  • Dag 詳細資訊 面板查看版本
使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...