Datenbewusstes Scheduling mit Assets

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Das Timing-Problem

 

  • Pipeline A lädt Rohdaten aus einer API
  • Pipeline B baut daraus ein Dashboard
  • Wenn beide mit cron laufen, kann B vor A fertig sein
  • Ein Zeitpuffer ist fragil

Das Timing-Problem

Data-Pipelines mit Airflow aufbauen

Was ist ein Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Eine Referenz auf Daten, identifiziert durch einen eindeutigen Namen
  • Eine URI kann angehängt werden, wenn das Asset eine konkrete Dateneinheit darstellt
  • Kann eine Datei, eine Datenbanktabelle oder beliebige Daten sein
Data-Pipelines mit Airflow aufbauen

Producer und Consumer

  • Tasks können Assets aktualisieren; dabei entsteht ein Asset-Event
  • Dags können bei Asset-Updates mit einem Asset-Schedule ausgelöst werden

 

Asset-Flow

Data-Pipelines mit Airflow aufbauen

Producer: Signalisieren mit Outlets

sales_data = Asset("s3://bucket/sales/daily.parquet")

@task(outlets=[sales_data]) def write_sales(): # Write data to S3 ...

 

  • Bei Erfolg protokolliert Airflow, dass das Asset aktualisiert wurde
  • Der Producer muss nichts über einen Consumer wissen
Data-Pipelines mit Airflow aufbauen

Consumer: auf ein Asset schedulen

sales_data = Asset("s3://bucket/sales/daily.parquet")

@dag(schedule=[sales_data]) def build_dashboard(): ...

 

  • Kein Cron-Ausdruck, nur die Asset-Referenz
  • Airflow triggert den Consumer automatisch nach dem Update
Data-Pipelines mit Airflow aufbauen

Bedingtes Scheduling

sales = Asset("s3://bucket/sales.parquet")
inventory = Asset("s3://bucket/inventory.parquet")

# Auf BEIDE Asset-Updates warten @dag(schedule=(sales & inventory)) def full_report(): ...
# Auslösen, wenn IRGENDEIN Asset updatet @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Data-Pipelines mit Airflow aufbauen

Assets mit der CLI prüfen

$ airflow assets list
name                                  | group    | uri
s3://data-lake/sales/daily.csv        | asset    | s3://data-lake/sales/daily.csv
$ airflow assets details --name "s3://data-lake/sales/daily.csv"
property_name    | property_value
=================+=====================                                         
name             | s3://data-lake/sales/daily.csv
                ...
updated_at       | 2026-04-14T10:00:50.820041Z
  • airflow assets list zeigt alle registrierten Assets
  • airflow assets details zeigt Metadaten inkl. updated_at
Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...