Transmiterea datelor între sarcini cu XCom

Introducere în Apache Airflow în Python

Mike Metzger

Data Engineer

Ce este XCom?

  • „Cross Communication"
    • Permite sarcinilor să comunice între ele
  • Stocat în baza de date de metadate Airflow
    • Transmite cantități mici de date
    • Nume de fișiere, URI, număr de rânduri

Ilustrație a XCom transmițând date mici între sarcini Airflow

Introducere în Apache Airflow în Python

Ce nu se trimite prin XCom

  • Fișiere mari
  • DataFrame-uri
  • Baze de date complete
  • Imagini mari

Ilustrație a tipurilor de date de evitat prin XCom, precum fișiere mari și DataFrame-uri

Introducere în Apache Airflow în Python

Implementarea XCom

  • Există mai multe moduri de a utiliza XCom
  • Ne concentrăm pe TaskFlow API
  • Extinde ceea ce am făcut deja cu @tasks
Introducere în Apache Airflow în Python

Exemplu 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()
Introducere în Apache Airflow în Python

Dependențe XCom

  • XCom definește automat ordinea dependențelor
  • Exemplu
    clean_data(get_data())
    
  • Echivalent conceptual cu get_data() >> clean_data()
  • Alt exemplu
     result = clean_data(get_data())
     result >> alert_when_complete()
    
Introducere în Apache Airflow în Python

Vizualizarea datelor XCom

Pagina XCom din Airflow cu valorile stocate după cheie, DAG și sarcină

Introducere în Apache Airflow în Python

Să exersăm!

Introducere în Apache Airflow în Python

Preparing Video For Download...