Data delen tussen taken

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Wat is XCom?

XCom push pull flow

 

  • XCom = communicatie tussen taken
  • Waarden staan in de metadatadatabase
  • Klassieke aanpak: expliciet xcom_push() en xcom_pull()
Data-pijplijnen bouwen met Airflow

XComArgs: de TaskFlow-methode

 

  • Een @task-returnwaarde wordt een XComArg
  • XComArg is een luie referentie, niet de echte data
  • Doorgeven aan een andere functie maakt de afhankelijkheid
  • Tijdens runtime lost Airflow het op via een pull uit XCom

XComArgs in TaskFlow

Data-pijplijnen bouwen met Airflow

XComArgs in code

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # afhankelijkheid + data
Data-pijplijnen bouwen met Airflow

Klassiek en TaskFlow combineren

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)
Data-pijplijnen bouwen met Airflow

XCom-beperkingen

XCom size limits by backend

 

  • Max. grootte hangt af van je database-backend
  • Moet JSON-serialiseerbaar zijn (dicts, lijsten, strings, getallen)
  • Elke waarde voegt belasting toe aan de database
Data-pijplijnen bouwen met Airflow

Aangepaste XCom-backends

Custom XCom backend routing

  • Routeer waarden naar een objectopslag
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: sla XComs alleen op in de objectopslag als de drempel is bereikt
Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...