1. Task Group 개념
- Task들의 모음
- Dag이 커져서 Task가 매우 많아지고 비슷한 구조가 많아지면 그룹화해서 좀더 관리하기 쉽도록 함
(꼭 써야 하는 것은 아니지만 유용하게 관리하는데 도움을 줄 수 있음) - 자세한 내용은 아래 링크에 있음
https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/dags.html#taskgroups
DAGs — Airflow Documentation
airflow.apache.org
- 정리하면
- Task Group 작성 방법은 2가지가 존재함
(데코레이터 & 클래스) - Task Group 안에 Task Group 중첩하여 정의 가능
- Task Group 간에도 Flow 정의 가능
- Group이 다르면 task_id가 같아도 무방
- Tooltip 파라미터를 이용해 ui 화면에서 task group에 대한 설명 제공 가능
(데코레이터 활용시 docstring으로도 가능)
- Task Group 작성 방법은 2가지가 존재함
2. Task Group 사용 예제코드
'''
dags_python_with_task_group.py
'''
from airflow import DAG
import pendulum
import datetime
from airflow.operators.python import PythonOperator
from airflow.decorators import task
from airflow.decorators import task_group
from airflow.utils.task_group import TaskGroup
with DAG(
dag_id='dags_python_with_task_group',
schedule=None,
start_date=pendulum.datetime(2023, 4, 1, tz='Asia/Seoul'),
catchup = False
) as dag:
def inner_func(**kwargs):
msg = kwargs.get('msg') or ''
print(msg)
# 방법1 task_group 데코레이터 사용
@task_group(group_id='first_group')
def group_1():
''' task_group 데코레이터를 이용한 첫 번째 그룹 '''
@task(task_id='inner_function1')
def inner_func1(**kwargs):
print('첫 번째 TaskGroup 내 첫 번째 task 입니다.')
inner_function2 = PythonOperator(
task_id = 'inner_function2',
python_callable=inner_func,
op_kwargs={'msg':'첫 번째 TaskGroup 내 두 번째 task입니다.'}
)
inner_func1() >> inner_function2
# 방법2 클래스 사용
with TaskGroup(group_id='second_group', tooltip='두 번째 그룹입니다') as group_2:
''' 여기에 적은 docstring은 표시되지 않습니다'''
@task(task_id='inner_function1')
def inner_func1(**kwargs):
print('두 번째 TaskGroup 내 첫 번째 task 입니다.')
inner_function2 = PythonOperator(
task_id = 'inner_function2',
python_callable=inner_func,
op_kwargs={'msg':'두 번째 TaskGroup 내 두 번째 task입니다.'}
)
inner_func1() >> inner_function2
group_1() >> group_2
실행결과 Graph


'Data 엔지니어링 > Airflow' 카테고리의 다른 글
| [Airflow] 지원되는 오퍼레이터 (0) | 2025.01.12 |
|---|---|
| [Airflow] Edge Label 사용하기 (0) | 2025.01.12 |
| [Airflow] Trigger Rule 설정하기 (0) | 2025.01.12 |
| [Airflow] BaseBranchOperator 로 분기처리하기 (0) | 2025.01.11 |
| [Airflow] @Task.branch로 분기처리하기 (0) | 2025.01.11 |