TaskFlow API

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

แนวทางแบบดั้งเดิม

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
การสร้าง Data Pipeline ด้วย Airflow

แนวทาง TaskFlow

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
การสร้าง Data Pipeline ด้วย Airflow

ทำไมต้องใช้ TaskFlow?

 

  • โค้ดกระชับขึ้น decorator แทน operator instance
  • กำหนด dependency อัตโนมัติ ด้วยค่าที่ส่งคืนระหว่าง task
  • อ่านง่าย DAG มีหน้าตาเหมือน Python script ทั่วไป
  • Classic operator ยังใช้ได้สำหรับ provider integration

TaskFlow API - ภาพรวม

การสร้าง Data Pipeline ด้วย Airflow

Airflow UI: มุมมอง Grid

มุมมอง Grid ใน Airflow

 

  • แต่ละคอลัมน์คือ DAG run หนึ่งครั้ง
  • แต่ละแถวคือ task หนึ่งรายการ
  • สีแสดงสถานะ: เขียว = สำเร็จ, แดง = ล้มเหลว
  • คลิกช่องเพื่อดู log และรายละเอียด
การสร้าง Data Pipeline ด้วย Airflow

Airflow UI: มุมมอง Graph

มุมมอง Graph ใน Airflow

 

  • แสดงภาพรวมของ DAG run เดียว
  • เห็น dependency และการรันแบบขนาน
  • ปรับรายละเอียดได้ในการตั้งค่า
การสร้าง Data Pipeline ด้วย Airflow

Airflow UI: การจัดการเวอร์ชัน DAG

ตัวบ่งชี้เวอร์ชัน DAG ใน Airflow UI

 

  • ติดตาม การเปลี่ยนแปลงโครงสร้าง โดยอัตโนมัติ
  • แต่ละ run เชื่อมกับ เวอร์ชันที่ใช้งานอยู่ในขณะนั้น
  • ดูเวอร์ชันได้ในแผง DAG details
การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...