Tổ chức Dag phức tạp với Task Group

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

Volker Janz

Senior Developer Advocate at Astronomer

Bài toán độ phức tạp

Dag phức tạp

  • Dag lớn khó điều hướng trong Graph view
  • Tên task lẫn vào nhau
  • Thành viên mới khó tìm phần mình phụ trách
Xây dựng Data Pipeline với 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 đặt định danh tùy chỉnh, mặc định là tên hàm
  • default_args áp dụng cho mọi task trong nhóm để tránh lặp cấu hình
  • Task group hiển thị như khối có thể thu/phóng trong UI
Xây dựng Data Pipeline với Airflow

Giao diện Airflow trông ra sao

Không dùng task group

Đồ thị đơn giản với ba task nối tiếp: extract_orders, transform_orders, load, cùng cấp

  • Mọi task ở cùng một cấp

Có task group

Đồ thị cho thấy khối process_orders đã thu gọn chứa extract và transform, tiếp theo là task load bên ngoài nhóm

  • Khối có thể thu/phóng trong UI
Xây dựng Data Pipeline với Airflow

Lồng nhóm và tên hiển thị tùy chỉnh

@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(),
    }
  • Có thể lồng task group cho pipeline phức tạp hơn
  • group_display_name đặt nhãn dễ đọc trong UI (hỗ trợ cả emoji)
Xây dựng Data Pipeline với Airflow

Mẫu factory

@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 là hàm Python có decorator
  • Bạn có thể gọi nhiều lần với tham số khác nhau
  • Mẫu factory giúp tái sử dụng logic pipeline
Xây dựng Data Pipeline với Airflow

Nguyên tắc nhóm

Nhóm theo miền hoặc mối quan tâm

  • Nhóm theo miền hoặc mối quan tâm, không theo loại operator
  • Dùng group_id cho tên lập trình rõ ràng, group_display_name cho nhãn UI
  • Dùng default_args để chia sẻ cấu hình như retries cho toàn bộ task trong nhóm
  • Áp dụng Định luật Miller, hơn 7 mục cấp cao có thể là dấu hiệu cần task group
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...