Operatory odroczone

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problem z sensorami

Sensors occupying worker slots

 

  • Tryb poke: sensor zajmuje slot przez cały czas
  • Długo działające sensory blokują inne zadania
Budowanie potoków danych z Airflow

Tryb reschedule

Tryb poke (domyślny)

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • Slot zajęty przez cały czas

Tryb reschedule

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • Slot zwalniany między sprawdzeniami
Budowanie potoków danych z Airflow

Tryb deferrable

Deferred sensor handled by the Triggerer

  • Zadanie przekazuje kontrolę procesowi Triggerer
  • Triggerer obsługuje setki równoległych oczekiwań dzięki async I/O
  • Podczas oczekiwania żaden slot nie jest zajęty
Budowanie potoków danych z Airflow

Włączanie trybu 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 nadal obowiązują
Budowanie potoków danych z Airflow

Kiedy używać każdego trybu

Tryb Zastosowanie Slot workera
poke Krótkie oczekiwania (< 1 min) Zajęty przez cały czas
reschedule Umiarkowane oczekiwania, bez Triggerer Zwalniany między sprawdzeniami
deferrable Długie oczekiwania, wiele sensorów Nigdy nie zajęty

 

  • Jeśli Triggerer działa, deferrable to niemal zawsze właściwy wybór
Budowanie potoków danych z Airflow

Globalne ustawienie deferrable

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • Wszystkie zgodne operatory automatycznie przechodzą w tryb deferrable
  • Kod DAG-a nie wymaga żadnych zmian
  • W razie potrzeby można wyłączyć dla konkretnego operatora: deferrable=False
Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...