Partager des données entre les tâches

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Qu'est-ce que XCom ?

Flux de push/pull XCom

 

  • XCom = échanges entre tâches
  • Valeurs stockées dans la base de métadonnées
  • Méthode classique : xcom_push() et xcom_pull() explicites
Créer des pipelines de données avec Airflow

XComArg : l'approche TaskFlow

 

  • La valeur de retour d'un @task devient un XComArg
  • XComArg est une référence paresseuse, pas la donnée elle‑même
  • La passer à une autre fonction crée la dépendance
  • À l'exécution, Airflow la résout en lisant dans XCom

XComArgs dans TaskFlow

Créer des pipelines de données avec Airflow

XComArg dans le code

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # dependency + data
Créer des pipelines de données avec Airflow

Mélanger classique et TaskFlow

get_timestamp = BashOperator(
    task_id="get_timestamp",
    bash_command="date +%Y-%m-%d")

@task
def log_timestamp(ts_value):
    print(f"Timestamp: {ts_value}")


log_timestamp(get_timestamp.output)
Créer des pipelines de données avec Airflow

Limites de XCom

Limites de taille XCom selon l'arrière-plan

 

  • La taille max dépend de votre moteur de base de données
  • Doit être sérialisable en JSON (dict, listes, chaînes, nombres)
  • Chaque valeur ajoute une charge à la base de données
Créer des pipelines de données avec Airflow

Arrière-plans XCom personnalisés

Acheminement vers un arrière-plan XCom personnalisé

  • Acheminer les valeurs vers un stockage d'objets
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD : n'enregistrer les XCom que dans le stockage d'objets si le seuil est atteint
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...