본문 바로가기
Data 엔지니어링/Airflow

[Airflow] Trigger Rule 설정하기

by imkus 2025. 1. 12.

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