상위: Airflow
요약
TaskGroup은 DAG 내에서 여러 Task를 논리적 단위로 그룹화하는 기능입니다. 복잡한 DAG의 가독성을 향상시키고, UI에서 그룹을 접고 펼칠 수 있습니다. 중첩 TaskGroup을 사용하여 계층 구조를 만들 수 있으며, 그룹 단위로 의존성을 설정할 수 있습니다.
TaskGroup이란?
기본 사용법
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 특징
| 항목 | 설명 |
|---|---|
| 목적 | 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"
다음 단계
- Task Dependencies - 그룹 간 의존성 설정
- DAG - DAG 구조 이해
- Task - Task 기본 개념