Xây dựng Data Pipeline với Airflow
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()
$$

$$

$$
# Không tự động chạy: kích hoạt thủ công (mặc định) @dag(schedule=None) def my_pipeline(): ...# Theo thời gian: chạy mỗi ngày lúc 6h sáng @dag(schedule="0 6 * * *") def daily_pipeline(): ...# Theo dữ liệu: chạy khi một Asset được cập nhật @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Classic operators
extract = PythonOperator(
task_id="extract",
python_callable=extract_fn)
extract >> transform
TaskFlow API
@task
def extract():
return {"users": 150}
data = extract()
transform(data)
$$

.py thậtpython3 filename.pydag.test() chạy toàn bộ Dag trong một tiến trìnhXây dựng Data Pipeline với Airflow