Deferrable オペレーター

Airflow によるデータパイプラインの構築

Volker Janz

Senior Developer Advocate at Astronomer

センサーの問題

ワーカースロットを占有するセンサー

 

  • Poke モード: センサーが常にワーカースロットを占有します
  • 長時間動作するセンサーは他のタスクを枯渇させます
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 が処理する Deferred センサー

  • タスクは Triggerer プロセスに処理を委譲します
  • Triggerer は非同期 I/O を使用して数百の待機を同時に処理します
  • 待機中にワーカースロットを消費しません
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 によるデータパイプラインの構築

各モードの使い分け

モード 適した用途 ワーカースロット
poke 短い待機(1 分未満) 常に占有
reschedule 中程度の待機、Triggerer 不使用時 チェック間に解放
deferrable 長い待機、センサー多数 使用しない

 

  • Triggerer が稼働中であれば、deferrable がほぼ常に最適なデフォルトです
Airflow によるデータパイプラインの構築

Deferrable のグローバル設定

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • 対応するすべてのオペレーターが自動的に Deferrable モードを使用します
  • DAG コードの変更は不要です
  • 必要に応じてオペレーターごとに deferrable=False で上書きできます
Airflow によるデータパイプラインの構築

練習しましょう!

Airflow によるデータパイプラインの構築

Preparing Video For Download...