Airflow ile Veri İş Hatları Oluşturma
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()
$$

$$

$$
# Otomatik çalıştırma yok: elle tetikle (varsayılan) @dag(schedule=None) def my_pipeline(): ...# Zaman tabanlı: her gün sabah 6'da çalışır @dag(schedule="0 6 * * *") def daily_pipeline(): ...# Veri farkındalıklı: bir Asset güncellenince çalışır @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Klasik operatörler
extract = PythonOperator(
task_id="extract",
python_callable=extract_fn)
extract >> transform
TaskFlow API
@task
def extract():
return {"users": 150}
data = extract()
transform(data)
$$

.py dosyalarını düzenlepython3 filename.py kullandag.test() tüm Dag'i tek işlemde çalıştırırAirflow ile Veri İş Hatları Oluşturma