Partajarea datelor între task-uri

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Ce este XCom?

Fluxul XCom push pull

 

  • XCom = comunicare între task-uri
  • Valorile sunt stocate în baza de date metadata
  • Abordare clasică: xcom_push() și xcom_pull() explicite
Construirea pipeline-urilor de date cu Airflow

XComArgs: metoda TaskFlow

 

  • Valoarea returnată de @task devine un XComArg
  • XComArg este o referință lazy, nu datele propriu-zise
  • Transmiterea lui altei funcții creează dependența
  • La execuție, Airflow îl rezolvă prin XCom

XComArgs în TaskFlow

Construirea pipeline-urilor de date cu Airflow

XComArgs în cod

@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
Construirea pipeline-urilor de date cu Airflow

Combinarea abordării clasice cu 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)
Construirea pipeline-urilor de date cu Airflow

Limitările XCom

Limite de dimensiune XCom pe backend

 

  • Dimensiunea maximă depinde de backend-ul bazei de date
  • Valorile trebuie să fie serializabile JSON (dicționare, liste, șiruri, numere)
  • Fiecare valoare adaugă sarcină bazei de date
Construirea pipeline-urilor de date cu Airflow

Backend-uri XCom personalizate

Rutare backend XCom personalizat

  • Rutează valorile către un object storage
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: stochează XCom-urile în object storage doar dacă pragul este atins
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...