Operatori differibili

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Il problema dei sensori

Sensori che occupano gli slot dei worker

 

  • Poke mode: il sensore occupa lo slot del worker per tutto il tempo
  • I sensori di lunga durata affamano le altre task
Creare data pipeline con Airflow

Reschedule mode

Poke mode (predefinito)

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • Slot occupato per tutto il tempo

Reschedule mode

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • Slot liberato tra i controlli
Creare data pipeline con Airflow

Modalità differibile

Sensore differito gestito da Triggerer

  • La task passa il controllo al processo Triggerer
  • Il Triggerer usa I/O asincrono per centinaia di attese concorrenti
  • Zero slot dei worker usati durante l'attesa
Creare data pipeline con Airflow

Abilitare la modalità differibile

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 e timeout si applicano comunque
Creare data pipeline con Airflow

Quando usare ogni modalità

Mode Ideale per Slot worker
poke Attese brevi (< 1 min) Occupato sempre
reschedule Attese medie, senza Triggerer Liberato tra i controlli
deferrable Attese lunghe, molti sensori Mai usato

 

  • Se il Triggerer è attivo, deferrable è quasi sempre la scelta giusta di default
Creare data pipeline con Airflow

Impostazione globale per differibile

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • Tutti gli operatori compatibili usano la modalità differibile in automatico
  • Nessuna modifica al codice del DAG necessaria
  • Puoi escludere per singolo operatore con deferrable=False se serve
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...