Data-Pipelines mit Airflow aufbauen
Volker Janz
Senior Developer Advocate at Astronomer

$$
$$

$$

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()
$$

$$

$$
# 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(): ...
Klassische Operatoren
extract = PythonOperator(
task_id="extract",
python_callable=extract_fn)
extract >> transform
TaskFlow API
@task
def extract():
return {"users": 150}
data = extract()
transform(data)
$$

.py-Dateien bearbeitenpython3 filename.pydag.test() führt den ganzen Dag in einem Prozess ausData-Pipelines mit Airflow aufbauen