TaskFlow API

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Cách tiếp cận cổ điển

def extract_data():
    return {"users": 150, "events": 4200}

with DAG("etl_pipeline") as dag: t1 = PythonOperator(task_id="extract", python_callable=extract_data) t2 = PythonOperator(task_id="summary", python_callable=print_summary)
t1 >> t2
Xây dựng Data Pipeline với Airflow

Cách tiếp cận TaskFlow

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
Xây dựng Data Pipeline với Airflow

Vì sao dùng TaskFlow?

 

  • Ít mã khuôn mẫu hơn, decorator thay cho operator instance
  • Phụ thuộc ngầm định, giá trị trả về tự nối các task
  • Dễ đọc, Dag như một script Python
  • Vẫn có operator cổ điển cho tích hợp với provider

TaskFlow API - visual

Xây dựng Data Pipeline với Airflow

Airflow UI: Grid view

Chế độ xem Grid của Airflow

 

  • Mỗi cột là một Dag run
  • Mỗi hàng là một task
  • Màu cho biết trạng thái: xanh = thành công, đỏ = thất bại
  • Nhấp vào một ô để xem log và chi tiết
Xây dựng Data Pipeline với Airflow

Airflow UI: Graph view

Chế độ xem Graph của Airflow

 

  • Tập trung vào một Dag run
  • Thể hiện phụ thuộc và thực thi song song
  • Điều chỉnh chi tiết trong phần cài đặt
Xây dựng Data Pipeline với Airflow

Airflow UI: Dag versioning

Chỉ báo phiên bản Dag trong Airflow UI

 

  • Tự động theo dõi thay đổi cấu trúc
  • Mỗi lần chạy liên kết tới phiên bản đang hoạt động khi đó
  • Tìm phiên bản trong bảng Dag details
Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...