Dela data mellan tasks

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Vad är XCom?

XCom push pull-flöde

 

  • XCom = kommunikation mellan tasks
  • Värden lagras i metadatadatabasen
  • Klassiskt sätt: explicit xcom_push() och xcom_pull()
Bygg datapipelines med Airflow

XComArgs: TaskFlow-sättet

 

  • Ett returvärde från @task blir ett XComArg
  • XComArg är en lazy-referens, inte det faktiska värdet
  • Att skicka det till en annan funktion skapar beroendet
  • Vid körning löser Airflow upp det via XCom

XComArgs i TaskFlow

Bygg datapipelines med Airflow

XComArgs i kod

@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
Bygg datapipelines med Airflow

Blanda klassisk och 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)
Bygg datapipelines med Airflow

XCom-begränsningar

XCom-storleksgränser per backend

 

  • Maxstorlek beror på databasbackend
  • Måste vara JSON-serialiserbart (dict, listor, strängar, tal)
  • Varje värde ökar belastningen på databasen
Bygg datapipelines med Airflow

Anpassade XCom-backends

Anpassad XCom-backend-routning

  • Dirigera värden till en objektlagring
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: lagrar XComs i objektlagringen först när tröskelvärdet nås
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...