Daten zwischen Tasks teilen

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Was ist XCom?

XCom push pull flow

 

  • XCom = Cross-Kommunikation zwischen Tasks
  • Werte liegen in der Metadaten-Datenbank
  • Klassisch: explizit xcom_push() und xcom_pull()
Data-Pipelines mit Airflow aufbauen

XComArgs: der TaskFlow-Weg

 

  • Ein @task-Rückgabewert wird zu einem XComArg
  • XComArg ist eine lazy reference, nicht die eigentlichen Daten
  • Beim Übergeben an eine andere Funktion entsteht die Abhängigkeit
  • Zur Laufzeit auflöst Airflow das, indem es aus XCom zieht

XComArgs in TaskFlow

Data-Pipelines mit Airflow aufbauen

XComArgs im Code

@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
Data-Pipelines mit Airflow aufbauen

Classic und TaskFlow mischen

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-Pipelines mit Airflow aufbauen

XCom-Einschränkungen

XCom size limits by backend

 

  • Maximalgröße hängt vom Datenbank-Backend ab
  • Muss JSON-serialisierbar sein (Dicts, Listen, Strings, Zahlen)
  • Jeder Wert erzeugt Last auf der Datenbank
Data-Pipelines mit Airflow aufbauen

Eigene XCom-Backends

Custom XCom backend routing

  • Werte an einen Object Storage routen
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: XComs nur im Object Storage speichern, wenn der Schwellwert erreicht ist
Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...