可延後的運算元

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

感測器的問題

感測器佔用 worker 插槽

 

  • Poke 模式:sensor 全程佔住 worker 插槽
  • 長時間執行的感測器會讓其他作業被餓死
使用 Airflow 建置資料管線

Reschedule 模式

Poke 模式(預設)

from airflow.providers.standard.sensors.filesystem \
    import FileSensor

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • 插槽全程被佔用

Reschedule 模式

from airflow.providers.standard.sensors.filesystem \
    import FileSensor

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • 檢查間插槽會釋放
使用 Airflow 建置資料管線

Deferrable 模式

Triggerer 處理延後的感測器

  • 作業交由 Triggerer 行程處理
  • Triggerer 以 非同步 I/O 同時等待數百項
  • 等待期間不佔用任何 worker 插槽
使用 Airflow 建置資料管線

啟用 deferrable

from airflow.providers.standard.sensors.filesystem import FileSensor

wait_for_data = FileSensor(
    task_id="wait_for_data",
    filepath="/data/incoming/report.csv",

deferrable=True,
poke_interval=30, timeout=3600, )

 

  • poke_intervaltimeout 仍然生效
使用 Airflow 建置資料管線

各模式何時用

模式 最適用 worker 插槽
poke 短等待(< 1 分鐘) 全程佔用
reschedule 中等等待,無 Triggerer 檢查間釋放
deferrable 長等待,感測器多 從不佔用

 

  • 若已啟動 Triggerer,deferrable 幾乎都是預設首選
使用 Airflow 建置資料管線

全域 deferrable 設定

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • 所有相容的運算元會自動使用 deferrable 模式
  • 不需修改 DAG 程式碼
  • 需要時可在單一運算元設 deferrable=False 覆寫
使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...