Sdílení dat mezi tasky

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Co je XCom?

Tok XCom push pull

 

  • XCom = komunikace mezi tasky
  • Hodnoty uložené v databázi metadat
  • Klasický přístup: explicitní xcom_push() a xcom_pull()
Tvorba datových pipeline s Airflow

XComArgs: způsob TaskFlow

 

  • Návratová hodnota @task se stane XComArg
  • XComArg je lazy reference, ne samotná data
  • Předání do jiné funkce vytvoří závislost
  • Za běhu Airflow hodnotu vyřeší načtením z XComu

XComArgs v TaskFlow

Tvorba datových pipeline s Airflow

XComArgs v kódu

@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
Tvorba datových pipeline s Airflow

Kombinace klasického přístupu a 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)
Tvorba datových pipeline s Airflow

Omezení XComu

Limity velikosti XCom podle backendu

 

  • Maximální velikost závisí na databázovém backendu
  • Musí být serializovatelné do JSON (dicts, lists, řetězce, čísla)
  • Každá hodnota přidává zátěž databázi
Tvorba datových pipeline s Airflow

Vlastní XCom backends

Směrování vlastního XCom backendu

  • Směrování hodnot do objektového úložiště
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: XComy se uloží do objektového úložiště až po dosažení prahu
Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...