Condivisione dei dati tra task

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Che cos'è XCom?

Flusso push/pull XCom

 

  • XCom = comunicazione tra task
  • Valori salvati nel database dei metadati
  • Approccio classico: xcom_push() e xcom_pull() espliciti
Creare data pipeline con Airflow

XComArg: l'approccio TaskFlow

 

  • Il valore di ritorno di un @task diventa un XComArg
  • XComArg è un riferimento lazy, non il dato reale
  • Passarlo a un'altra funzione crea la dipendenza
  • In runtime, Airflow lo risolve estraendolo da XCom

XComArg in TaskFlow

Creare data pipeline con Airflow

XComArg nel codice

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

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


config = get_config() # XComArg
run_pipeline(config) # dipendenza + dati
Creare data pipeline con Airflow

Mescolare classico e 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)
Creare data pipeline con Airflow

Limitazioni di XCom

Limiti dimensione XCom per backend

 

  • La dimensione massima dipende dal database backend
  • Devono essere serializzabili in JSON (dict, liste, stringhe, numeri)
  • Ogni valore aggiunge carico al database
Creare data pipeline con Airflow

Backend XCom personalizzati

Instradamento verso backend XCom personalizzato

  • Instrada i valori a uno object storage
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: salva gli XCom nello object storage solo se si raggiunge la soglia
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...