1. Trigger Rule 종류

- 위와 같이 Task 관계가 구성된 경우 상위 Task (task1, 2, 3)이 모두 성공시에만 task4가 수행됨
- Trigger Rule을 이용하면 상위 Task들의 상태에 따라 수행여부를 결정할 수 있음
- Trigger Rule 종류
| all_success (기본값) | 상위 task가 모두 성공하면 실행 |
| all_failed | 상위 task가 모두 실패하면 실행 |
| all_done | 상위 task가 모두 수행되면 실행 (실패도 수행된것에 포함) |
| all_skipped | 상위 task가 모두 skipped 상태면 실행 |
| one_failed | 상위 task 중 하나 이상 실패하면 실행 (모든 상위 task 완료를 기다리지 않음) |
| one_success | 상위 task 중 하나 이상 성공하면 실행 (모든 상위 task 완료를 기다리지 않음) |
| one_done | 상위 task 중 하나 이상 성공 또는 실패 하면 실행 |
| none_failed | 상위 task 중 실패가 없는 경우 실행 (성공 또는 skipped 상태) |
| none_failed_min_one_success | 상위 task 중 실패가 없고 성공한 task가 적어도 1개 이상이면 실행 |
| none_skpped | skip 된 상위 task가 없으면 실행 (상위 task가 성공, 실패여도 무방) |
| always | 언제나 실행 |
2. Trigger Rule (all_done) 사용 예제 코드
'''
dags_python_with_trigger_rule_eg1.py
'''
from airflow import DAG
from airflow.decorators import task
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
from airflow.exceptions import AirflowException
import pendulum
with DAG(
dag_id='dag_python_with_trigger_rule_eg1',
start_date=pendulum.datetime(2023, 4, 1, tz='Asia/Seoul'),
schedule=None,
catchup=False
) as dag:
bash_upstream_1 = BashOperator(
task_id = 'bash_upstream_1',
bash_command = 'echo upstream1'
)
@task(task_id='python_upstream_1')
def python_upstream_1():
raise AirflowException('downstream_1 Exception!') # 해당 task가 실패 하도록 함
@task(task_id='python_upstream_2')
def python_upstream_2():
print('정상 처리')
@task(task_id='python_downstream_1', trigger_rule='all_done') # trigger_rule에 all_done을 적용하여 상위 task가 모두 종료되어야 실행됨
def python_downstream_1():
print('정상 처리')
[bash_upstream_1, python_upstream_1(), python_upstream_2()] >> python_downstream_1()
실행결과 Graph

2. Trigger Rule (none_skipped) 사용 예제 코드
'''
dags_python_with_trigger_rule_eg2.py
'''
from airflow import DAG
from airflow.decorators import task
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
from airflow.exceptions import AirflowException
import pendulum
with DAG(
dag_id='dags_python_with_trigger_rule_eg2',
start_date=pendulum.datetime(2023,4,1, tz='Asia/Seoul'),
schedule=None,
catchup=False
) as dag:
@task.branch(task_id='branching')
def random_branch():
import random
item_lst = ['A', 'B', 'C']
selected_item = random.choice(item_lst)
if selected_item == 'A':
return 'task_a'
elif selected_item == 'B':
return 'task_b'
elif selected_item == 'C':
return 'task_c'
task_a = BashOperator(
task_id='task_a',
bash_command='echo upstram1'
)
@task(task_id='task_b')
def task_b():
print('정상 처리')
@task(task_id='task_c')
def task_c():
print('정상 처리')
@task(task_id='task_d', trigger_rule='none_skipped')
def task_d():
print('정상 처리')
random_branch() >> [task_a, task_b(), task_c()] >> task_d()
실행결과 Graph

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