Współdzielenie danych między zadaniami

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Czym jest XCom?

XCom push pull flow

 

  • XCom = komunikacja między zadaniami (cross-communication)
  • Wartości przechowywane w bazie metadanych
  • Klasyczne podejście: jawne xcom_push() i xcom_pull()
Budowanie potoków danych z Airflow

XComArgs: podejście TaskFlow

 

  • Wartość zwracana przez @task staje się XComArg
  • XComArg to leniwa referencja, a nie same dane
  • Przekazanie jej do innej funkcji tworzy zależność
  • W czasie wykonania Airflow rozwiązuje ją, pobierając dane z XCom

XComArgs in TaskFlow

Budowanie potoków danych z Airflow

XComArgs w kodzie

@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
Budowanie potoków danych z Airflow

Mieszanie klasycznego podejścia z 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)
Budowanie potoków danych z Airflow

Ograniczenia XCom

XCom size limits by backend

 

  • Maksymalny rozmiar zależy od używanego backendu bazy danych
  • Dane muszą być serializowalne do JSON (słowniki, listy, ciągi znaków, liczby)
  • Każda wartość zwiększa obciążenie bazy danych
Budowanie potoków danych z Airflow

Niestandardowe backendy XCom

Custom XCom backend routing

  • Kieruj wartości do magazynu obiektów
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: dane trafiają do magazynu obiektów dopiero po przekroczeniu progu
Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...