Transférer des données entre tâches avec XCom

Introduction à Apache Airflow en Python

Mike Metzger

Data Engineer

Qu'est-ce que XCom ?

  • « Cross Communication »
    • Permet aux tâches de se parler
  • Stocké dans la base de métadonnées d'Airflow
    • Transmet de petites quantités de données
    • Noms de fichiers, URI, dénombrements de lignes

Illustration du transfert de petites données entre tâches Airflow avec XCom

Introduction à Apache Airflow en Python

Ce qu'il ne faut pas envoyer via XCom

  • Gros fichiers
  • DataFrames
  • Bases de données complètes
  • Grandes images

Illustration des types de données à éviter d'envoyer via XCom, comme les gros fichiers et les DataFrames

Introduction à Apache Airflow en Python

Implémenter XCom

  • Plusieurs façons d'utiliser XCom
  • Nous nous concentrerons sur l'API TaskFlow
  • Prolonge ce que nous avons déjà fait avec @tasks
Introduction à Apache Airflow en Python

Exemple XCom

@dag(dag_id='Example_XCom')
def example_xcom():

@task def get_data(): return data
@task(multiple_outputs=True) def clean_data(sourcedata): return clean(sourcedata) # Example, not implemented
clean_data(get_data()) example_xcom()
Introduction à Apache Airflow en Python

Dépendances XCom

  • Les XCom définissent automatiquement l'ordre des dépendances
  • Exemple
    clean_data(get_data())
    
  • Conceptuellement équivalent à get_data() >> clean_data()
  • Autre exemple
     result = clean_data(get_data())
     result >> alert_when_complete()
    
Introduction à Apache Airflow en Python

Afficher les données XCom

Page XCom d'Airflow listant les valeurs stockées par clé, DAG et tâche

Introduction à Apache Airflow en Python

Passons à la pratique !

Introduction à Apache Airflow en Python

Preparing Video For Download...