Airflow: Grundkonzepte

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Lern deinen Instructor kennen

  Profilfoto von Volker Janz

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ Jahre Data Engineer im Gaming
  • Arbeitet mit Airflow seit Version 1.x
  • Speaker, Mentor und Newsletter-Lead bei Data Engineer Things
Data-Pipelines mit Airflow aufbauen

Was du bauen wirst

 

  • Schreibe Dags mit der TaskFlow API
  • Baue dynamische Workflows mit Task Mapping und Asset-basiertem Scheduling
  • Geh mit Retries und Callbacks robust mit Fehlern um
  • Führe SQL-Workloads mit Airflow aus

Visualisierung der Inhalte je Kapitel

Data-Pipelines mit Airflow aufbauen

Bevor wir starten

 

$$

  • Sicher im Umgang mit Dags, Tasks und Operatoren
  • Vertraut mit Scheduling-Grundlagen

Introduction to Airflow - Kursseite

Data-Pipelines mit Airflow aufbauen

Kurzer Auffrischer

from airflow.sdk import dag, task

@dag
def star_wars_dag():


@task def get_star_wars_person(): import requests return requests.get("https://swapi.dev/api/people/1/").json()
@task.bash def print_name(person): return f"echo '{person['name']}'"
person = get_star_wars_person() print_name(person) star_wars_dag()
  • Ein Dag ist eine Sammlung von Tasks mit Abhängigkeiten
  • Tasks sind einzelne Arbeitsschritte
  • Operatoren/Decorator legen fest, was jede Task tut
  • Abhängigkeiten bestimmen die Ausführungsreihenfolge

$$

Einfacher Dag

Data-Pipelines mit Airflow aufbauen

Airflow-Architektur

$$

Airflow-3-Architektur

 

$$

  • Orchestrierung: Scheduler, Dag Processor
  • Ausführung: Worker, Triggerer
  • Interface & Storage: API Server, Metadata DB
Data-Pipelines mit Airflow aufbauen

Scheduling-Ansätze

# Keine automatischen Läufe: manuell triggern (Standard)
@dag(schedule=None)
def my_pipeline(): ...

# Zeitbasiert: läuft täglich um 6 Uhr @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Datenbasiert: läuft, wenn ein Asset aktualisiert wird @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Data-Pipelines mit Airflow aufbauen

Zwei Wege, Dags zu schreiben

Klassische Operatoren

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Wähle dies, wenn keine Decorators verfügbar sind

TaskFlow API

@task
def extract():
    return {"users": 150}

data = extract()
transform(data)
  • Einfache Python-Decorator
  • Weniger Boilerplate
  • Kombinierbar mit klassischen Operatoren
Data-Pipelines mit Airflow aufbauen

Übungen in diesem Kurs

$$

IDE-Übungs-Screenshot

 

  • IDE-Übungen: echte .py-Dateien bearbeiten
  • Klick auf "Run this file" oder nutze python3 filename.py
  • dag.test() führt den ganzen Dag in einem Prozess aus
Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...