การจัดระเบียบ DAG ที่ซับซ้อนด้วย Task Groups

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

Volker Janz

Senior Developer Advocate at Astronomer

ปัญหาความซับซ้อน

Complex Dag

  • DAG ขนาดใหญ่ นำทางยาก ใน Graph view
  • ชื่อ task ปนกันจนสับสน
  • สมาชิกใหม่ในทีมหา task ของตัวเองได้ยาก
การสร้าง Data Pipeline ด้วย Airflow

@task_group

from airflow.sdk import dag, task, task_group

@task_group(
    group_id="ingest_orders",
    default_args={"retries": 3},
)

def process_orders(): @task def extract_orders(): return [{"id": 1, "amount": 99.99}] @task def transform_orders(orders): return [{"id": o["id"], "total": o["amount"] * 1.08} for o in orders] return transform_orders(extract_orders())
  • group_id กำหนด identifier ที่ต้องการ หากไม่ระบุจะใช้ชื่อฟังก์ชัน
  • default_args ใช้กับ ทุก task ในกลุ่ม ช่วยลดโค้ดซ้ำ
  • Task group แสดงเป็นบล็อกที่ขยาย/ย่อได้ใน UI
การสร้าง Data Pipeline ด้วย Airflow

มุมมองใน Airflow UI

ไม่มี task group

กราฟแบบง่ายแสดง task สามขั้นตอนเรียงกัน ได้แก่ extract_orders, transform_orders และ load อยู่ในระดับเดียวกัน

  • ทุก task อยู่ในระดับเดียวกัน

มี task group

กราฟแสดงบล็อกที่ย่ออยู่ชื่อ process_orders ซึ่งประกอบด้วย extract และ transform ตามด้วย task ชื่อ load ที่อยู่นอกกลุ่ม

  • บล็อกที่ขยาย/ย่อได้ใน UI
การสร้าง Data Pipeline ด้วย Airflow

การซ้อนกลุ่มและชื่อแสดงผลแบบกำหนดเอง

@task_group(group_display_name="Process All")
def process_all():

    @task_group(group_display_name="Ingest Orders")
    def orders():
        return transform(extract())

    @task_group(group_display_name="Process Returns")
    def returns():
        return transform(extract())

    return {
        "orders": orders(),
        "returns": returns(),
    }
  • Task group ซ้อนกันได้สำหรับ pipeline ที่ซับซ้อนขึ้น
  • group_display_name กำหนด ชื่อที่อ่านได้ ใน UI (รองรับ emoji ด้วย)
การสร้าง Data Pipeline ด้วย Airflow

Factory pattern

@task_group
def process_source(source_name, source_path):

    @task
    def extract():
        return read_data(source_path)

    @task
    def transform(data):
        return clean_data(data)

    return transform(extract())

# Reuse the same pattern for different sources orders = process_source("orders", "/data/orders.csv") returns = process_source("returns", "/data/returns.csv") events = process_source("events", "/data/events.csv")
  • Task group คือ ฟังก์ชัน Python ที่ใช้ decorator
  • เรียกใช้ ได้หลายครั้ง ด้วย parameter ที่ต่างกัน
  • Factory pattern นี้ช่วย นำ logic ของ pipeline กลับมาใช้ซ้ำ
การสร้าง Data Pipeline ด้วย Airflow

แนวทางการจัดกลุ่ม

จัดกลุ่มตาม domain หรือ concern

  • จัดกลุ่มตาม domain หรือ concern ไม่ใช่ตามประเภท operator
  • ใช้ group_id สำหรับ ชื่อในโค้ด และ group_display_name สำหรับ ชื่อใน UI
  • ใช้ default_args เพื่อแชร์การตั้งค่า เช่น retries ให้ทุก task ในกลุ่ม
  • ใช้ Miller's Law เป็นแนวทาง — หาก มีรายการระดับบนสุดเกิน 7 รายการ อาจถึงเวลาแบ่ง task group
การสร้าง Data Pipeline ด้วย Airflow

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

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

Preparing Video For Download...