Opérateurs reportables

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Le problème des capteurs

Des capteurs occupant des emplacements de travailleurs

 

  • Mode poke : le capteur garde l'emplacement du travailleur tout du long
  • Les capteurs de longue durée affament les autres tâches
Créer des pipelines de données avec Airflow

Mode reschedule

Mode poke (par défaut)

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • Emplacement occupé tout du long

Mode reschedule

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

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • Emplacement libéré entre les vérifications
Créer des pipelines de données avec Airflow

Mode reportable

Capteur différé géré par le Triggerer

  • La tâche remet l'attente au processus Triggerer
  • Le Triggerer utilise des E/S asynchrones pour des centaines d'attentes concurrentes
  • Zéro emplacement de travailleur utilisé pendant l'attente
Créer des pipelines de données avec Airflow

Activer le mode reportable

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 et timeout s'appliquent toujours
Créer des pipelines de données avec Airflow

Quand utiliser chaque mode

Mode Idéal pour Emplacement de travailleur
poke Courtes attentes (< 1 min) Retenu tout du long
reschedule Attentes moyennes, sans Triggerer Libéré entre vérifications
reportable Longues attentes, nombreux capteurs Jamais utilisé

 

  • Si le Triggerer fonctionne, le mode reportable est presque toujours le bon choix par défaut
Créer des pipelines de données avec Airflow

Paramètre global du mode reportable

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • Tous les opérateurs compatibles utilisent le mode reportable automatiquement
  • Aucun changement requis dans le code du DAG
  • Outrepasser par opérateur avec deferrable=False au besoin
Créer des pipelines de données avec Airflow

Passons à la pratique !

Créer des pipelines de données avec Airflow

Preparing Video For Download...