전체 그래프
Airflow

TaskGroup

data-engineeringtool/airflowtaskgroup

상위: Airflow

요약

TaskGroup은 DAG 내에서 여러 Task를 논리적 단위로 그룹화하는 기능입니다. 복잡한 DAG의 가독성을 향상시키고, UI에서 그룹을 접고 펼칠 수 있습니다. 중첩 TaskGroup을 사용하여 계층 구조를 만들 수 있으며, 그룹 단위로 의존성을 설정할 수 있습니다.

TaskGroup이란?

  • DAG 내에서 여러 Task논리적 단위로 그룹화
  • DAG의 복잡성 감소, 유지보수성 향상
  • UI에서 그룹 단위로 접고 펼치기 가능
  • 성능 저하 없음 (DAG 내에서 실행)

기본 사용법

from airflow.utils.task_group import TaskGroup
from airflow.operators.bash import BashOperator

with DAG('example_dag', ...) as dag:
    
    start = EmptyOperator(task_id='start')
    
    with TaskGroup("processing_group") as processing:
        task_1 = BashOperator(task_id='task_1', bash_command='echo 1')
        task_2 = BashOperator(task_id='task_2', bash_command='echo 2')
        task_3 = BashOperator(task_id='task_3', bash_command='echo 3')
        
        task_1 >> task_2 >> task_3
    
    end = EmptyOperator(task_id='end')
    
    start >> processing >> end

중첩 TaskGroup

with DAG("nested_taskgroup_example", ...) as dag:

    start = EmptyOperator(task_id="start")

    with TaskGroup("group_1") as group_1:

        task_1 = BashOperator(
            task_id="task_1",
            bash_command="echo Task 1 실행"
        )

        # 내부 TaskGroup
        with TaskGroup("inner_group_1") as inner_group_1:
            inner_task = BashOperator(
                task_id="inner_task",
                bash_command="echo Inner Group 실행"
            )

        task_2 = BashOperator(task_id="task_2", bash_command="echo Task 2")
        task_3 = BashOperator(task_id="task_3", bash_command="echo Task 3")
        task_4 = BashOperator(task_id="task_4", bash_command="echo Task 4")

        task_1 >> inner_group_1 >> [task_2, task_3] >> task_4

    end = EmptyOperator(task_id="end")

    start >> group_1 >> end

TaskGroup

TaskGroup

TaskGroup 특징

항목설명
목적Task를 논리적으로 그룹화하여 DAG 가독성 개선
성능DAG 내에서 실행되므로 성능 저하 없음
병렬 실행병렬 실행 가능
UI 표현UI에서 그룹을 접거나 펼칠 수 있음
의존성 관리DAG 내에서 Task 간 의존성을 쉽게 설정 가능
실행 방식단순한 Task 묶음으로 실행 관리

TaskGroup ID

자동 생성되는 Task ID

with TaskGroup("group_1") as group:
    task = BashOperator(task_id="task", ...)

# 실제 task_id: "group_1.task"

중첩 그룹의 Task ID

with TaskGroup("outer") as outer:
    with TaskGroup("inner") as inner:
        task = BashOperator(task_id="task", ...)

# 실제 task_id: "outer.inner.task"

실전 패턴

ETL 파이프라인

with DAG('etl_pipeline', ...) as dag:
    
    start = EmptyOperator(task_id='start')
    
    with TaskGroup("extract") as extract_group:
        extract_db = PythonOperator(task_id='from_db', ...)
        extract_api = PythonOperator(task_id='from_api', ...)
        extract_file = PythonOperator(task_id='from_file', ...)
    
    with TaskGroup("transform") as transform_group:
        clean = PythonOperator(task_id='clean', ...)
        enrich = PythonOperator(task_id='enrich', ...)
        aggregate = PythonOperator(task_id='aggregate', ...)
        
        clean >> enrich >> aggregate
    
    with TaskGroup("load") as load_group:
        load_warehouse = PythonOperator(task_id='to_warehouse', ...)
        load_cache = PythonOperator(task_id='to_cache', ...)
    
    start >> extract_group >> transform_group >> load_group

데이터 소스별 그룹화

sources = ['mysql', 'postgres', 'mongodb']

for source in sources:
    with TaskGroup(f"process_{source}") as group:
        extract = PythonOperator(task_id='extract', ...)
        transform = PythonOperator(task_id='transform', ...)
        load = PythonOperator(task_id='load', ...)
        
        extract >> transform >> load

TaskGroup vs SubDAG

TaskGroup (권장)

with TaskGroup("group") as group:
    task1 >> task2

장점:

  • 간단한 구문
  • DAG 내에서 직접 실행
  • 성능 오버헤드 없음
  • UI에서 접고 펼치기

SubDAG (Deprecated)

from airflow.operators.subdag import SubDagOperator

subdag_task = SubDagOperator(
    task_id='subdag',
    subdag=create_subdag(...)
)

단점:

  • 복잡한 구성
  • 별도 DAG Run 생성 (성능 저하)
  • Airflow 2.0부터 TaskGroup 권장

주의사항

TaskGroup ID는 유일해야 함

#  같은 ID 사용 불가
with TaskGroup("group") as group1:
    ...

with TaskGroup("group") as group2:  # 에러!
    ...

#  다른 ID 사용
with TaskGroup("group_1") as group1:
    ...

with TaskGroup("group_2") as group2:
    ...

Task ID 충돌 주의

with TaskGroup("group") as group:
    task = BashOperator(task_id="task", ...)
    # 실제 ID: "group.task"

# 같은 이름 사용 가능 (다른 그룹)
with TaskGroup("other_group") as other:
    task = BashOperator(task_id="task", ...)
    # 실제 ID: "other_group.task"

다음 단계

참고 자료