Operatori deferrabili

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problema senzorilor

Senzori care ocupă sloturi de worker

 

  • Modul poke: senzorul ocupă slotul de worker tot timpul
  • Senzorii cu rulare lungă blochează alte sarcini
Construirea pipeline-urilor de date cu Airflow

Modul reschedule

Modul poke (implicit)

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • Slot ocupat tot timpul

Modul reschedule

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • Slot eliberat între verificări
Construirea pipeline-urilor de date cu Airflow

Modul deferrable

Senzor deferrabil gestionat de Triggerer

  • Sarcina transferă controlul procesului Triggerer
  • Triggerer folosește async I/O pentru sute de așteptări simultane
  • Zero sloturi de worker utilizate în timpul așteptării
Construirea pipeline-urilor de date cu Airflow

Activarea modului 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_interval și timeout se aplică în continuare
Construirea pipeline-urilor de date cu Airflow

Când să folosești fiecare mod

Mod Recomandat pentru Slot de worker
poke Așteptări scurte (< 1 min) Ocupat tot timpul
reschedule Așteptări moderate, fără Triggerer Eliberat între verificări
deferrable Așteptări lungi, mulți senzori Neutilizat

 

  • Dacă Triggerer rulează, deferrable este aproape întotdeauna opțiunea corectă
Construirea pipeline-urilor de date cu Airflow

Setare globală deferrable

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • Toți operatorii compatibili folosesc modul deferrable automat
  • Nu sunt necesare modificări în codul DAG
  • Suprascrie per operator cu deferrable=False dacă este nevoie
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...