Przekazywanie danych między zadaniami za pomocą XCom

Wprowadzenie do Apache Airflow w Pythonie

Mike Metzger

Data Engineer

Czym jest XCom?

  • „Cross Communication"
    • Umożliwia komunikację między zadaniami
  • Przechowywane w bazie metadanych Airflow
    • Przekazuje małe ilości danych
    • Nazwy plików, URI, liczby wierszy

Ilustracja przedstawiająca przekazywanie małych danych między zadaniami Airflow za pomocą XCom

Wprowadzenie do Apache Airflow w Pythonie

Czego nie przesyłać przez XCom

  • Duże pliki
  • DataFrames
  • Pełne bazy danych
  • Duże obrazy

Ilustracja typów danych, których nie należy przesyłać przez XCom, np. dużych plików i DataFrames

Wprowadzenie do Apache Airflow w Pythonie

Implementacja XCom

  • XCom można implementować na wiele sposobów
  • Skupimy się na TaskFlow API
  • Rozszerzenie dotychczasowego użycia @task
Wprowadzenie do Apache Airflow w Pythonie

Przykład 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()
Wprowadzenie do Apache Airflow w Pythonie

Zależności XCom

  • XCom automatycznie definiuje kolejność zależności
  • Przykład
    clean_data(get_data())
    
  • Odpowiednik get_data() >> clean_data()
  • Kolejny przykład
     result = clean_data(get_data())
     result >> alert_when_complete()
    
Wprowadzenie do Apache Airflow w Pythonie

Przeglądanie danych XCom

Strona XCom w Airflow z listą przechowywanych wartości według klucza, DAG-a i zadania

Wprowadzenie do Apache Airflow w Pythonie

Czas na praktykę!

Wprowadzenie do Apache Airflow w Pythonie

Preparing Video For Download...