Deferrable operators

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Sensorproblemet

Sensorer som upptar worker-platser

 

  • Poke-läge: sensorn håller worker-platsen hela tiden
  • Långvariga sensorer svälter ut andra uppgifter
Bygg datapipelines med Airflow

Reschedule-läge

Poke-läge (standard)

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • Platsen upptas hela tiden

Reschedule-läge

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • Platsen frigörs mellan kontroller
Bygg datapipelines med Airflow

Deferrable-läge

Uppskjuten sensor hanteras av Triggerer

  • Uppgiften lämnar över till Triggerer-processen
  • Triggerer använder asynkron I/O för hundratals samtidiga väntetillstånd
  • Inga worker-platser används under väntan
Bygg datapipelines med Airflow

Aktivera 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 och timeout gäller fortfarande
Bygg datapipelines med Airflow

När ska vilket läge användas?

Läge Passar för Worker-plats
poke Korta väntetider (< 1 min) Upptas hela tiden
reschedule Måttliga väntetider, ingen Triggerer Frigörs mellan kontroller
deferrable Långa väntetider, många sensorer Används aldrig

 

  • Om Triggerer körs är deferrable nästan alltid rätt val
Bygg datapipelines med Airflow

Global deferrable-inställning

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • Alla kompatibla operatorer använder deferrable-läge automatiskt
  • Inga ändringar i DAG-koden behövs
  • Åsidosätt per operator med deferrable=False vid behov
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...