DAG (DAGs)

DAG (DAGs)

DAG는 워크플로우를 실행하는 데 필요한 모든 것을 담아내는 모델이에요. DAG가 가지는 주요 속성들을 먼저 훑어볼게요.

  • Schedule (스케줄): 워크플로우를 언제 실행할지.
  • Tasks (태스크): tasks는 워커(worker)에서 실행되는 작업의 단위예요.
  • Task Dependencies (태스크 의존성): tasks가 어떤 순서와 조건에서 실행될지.
  • Callbacks (콜백): 워크플로우 전체가 완료되었을 때 취할 동작.
  • Additional Parameters (추가 파라미터): 그 외에도 많은 운영 관련 세부 사항들.

기본적인 예시 DAG 하나를 볼게요.

../_images/basic_dag.png

이 DAG는 A, B, C, D 네 개의 태스크를 정의하고, 이들이 실행되는 순서와 서로 간의 의존 관계를 정해요. 그리고 DAG를 얼마나 자주 실행할지도 정하는데, 예를 들어 "내일부터 5분마다" 또는 "2020년 1월 1일부터 매일" 같은 방식이지요.

DAG 자체는 태스크 내부에서 무엇이 일어나는지에는 관심이 없어요. DAG가 신경 쓰는 것은 태스크를 어떻게 실행할지, 즉 실행 순서, 재시도 횟수, 타임아웃 여부 등이에요.

Note

"DAG"라는 용어는 수학 개념인 "directed acyclic graph(방향성 비순환 그래프)"에서 왔지만, Airflow에서의 의미는 수학적 DAG 개념의 문자 그대로의 데이터 구조를 훨씬 넘어서 진화했어요. 그래서 Airflow에서는 'Dag'라는 용어를 쓰기로 결정했어요.

DAG 선언하기 (Declaring a Dag)

DAG를 선언하는 방법은 세 가지가 있어요. 먼저 with 문(컨텍스트 매니저)을 쓰는 방법인데, 블록 안에 들어 있는 모든 것을 DAG에 암묵적으로 추가해 줘요.

import datetime

from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator

with DAG(
    dag_id="my_dag_name",
    start_date=datetime.datetime(2021, 1, 1),
    schedule="@daily",
):
    EmptyOperator(task_id="task")

아니면 표준 생성자(constructor)를 사용해, 사용하는 모든 operator에 DAG를 넘겨줄 수도 있어요.

import datetime

from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator

my_dag = DAG(
    dag_id="my_dag_name",
    start_date=datetime.datetime(2021, 1, 1),
    schedule="@daily",
)
EmptyOperator(task_id="task", dag=my_dag)

아니면 @dag 데코레이터를 써서 함수를 DAG 생성기(generator)로 만들 수도 있어요.

import datetime

from airflow.sdk import dag
from airflow.providers.standard.operators.empty import EmptyOperator


@dag(start_date=datetime.datetime(2021, 1, 1), schedule="@daily")
def generate_dag():
    EmptyOperator(task_id="task")


generate_dag()

DAG는 실행할 Tasks가 없으면 아무것도 아니에요. 그리고 태스크는 보통 Operators, Sensors, TaskFlow 중 하나의 형태로 오게 돼요.

태스크 의존성 (Task Dependencies)

하나의 Task/Operator는 보통 혼자 존재하지 않아요. 다른 태스크에 대한 의존성(그 태스크의 upstream, 즉 상위에 있는 것들)을 갖고, 또 다른 태스크가 그 태스크에 의존하기도 해요(그 태스크의 downstream, 즉 하위에 있는 것들). 태스크 사이의 이런 의존성을 선언하는 것이 바로 DAG 구조를 만드는 일이에요.

개별 태스크 의존성을 선언하는 방법은 크게 두 가지가 있어요. 권장하는 방법은 >><< 연산자를 쓰는 것이에요.

first_task >> [second_task, third_task]
third_task << fourth_task

아니면 좀 더 명시적인 set_upstreamset_downstream 메서드를 쓸 수도 있어요.

first_task.set_downstream([second_task, third_task])
third_task.set_upstream(fourth_task)

더 복잡한 의존성을 선언하는 단축 방법도 있어요. 태스크 리스트가 다른 태스크 리스트에 의존하게 만들고 싶다면 위 두 가지 방법 모두 쓸 수 없어서, cross_downstream 을 써야 해요.

from airflow.sdk import cross_downstream

# 아래 코드를 대체해요
# [op1, op2] >> op3
# [op1, op2] >> op4
cross_downstream([op1, op2], [op3, op4])

그리고 의존성을 연쇄적으로(chaining) 연결하고 싶다면 chain 을 쓸 수 있어요.

from airflow.sdk import chain

# op1 >> op2 >> op3 >> op4 를 대체해요
chain(op1, op2, op3, op4)

# 동적으로도 할 수 있어요
chain(*[EmptyOperator(task_id=f"op{i}") for i in range(1, 6)])

chain 은 같은 크기의 리스트에 대해 쌍(pairwise) 의존성도 만들 수 있어요. (이것은 cross_downstream 이 만드는 cross 의존성과는 달라요!)

from airflow.sdk import chain

# 아래 코드를 대체해요
# op1 >> op2 >> op4 >> op6
# op1 >> op3 >> op5 >> op6
chain(op1, [op2, op3], [op4, op5], op6)

DAG 로딩 (Loading Dags)

Airflow는 DAG 번들(bundle) 안의 Python 소스 파일에서 DAG를 로딩해요. 각 파일을 가져와서 실행한 뒤, 그 파일에서 DAG 객체들을 로드하지요.

즉, Python 파일 하나에 여러 개의 DAG를 정의할 수도 있고, 아주 복잡한 DAG 하나를 import를 사용해 여러 Python 파일에 걸쳐 나눠 담을 수도 있어요.

단, Airflow가 Python 파일에서 DAG를 로드할 때는 최상위 레벨(top level) 의 객체 중 DAG 인스턴스인 것들만 가져온다는 점을 유의해야 해요. 예를 들어 이런 DAG 파일을 살펴볼게요.

dag_1 = DAG('this_dag_will_be_discovered')

def my_function():
    dag_2 = DAG('but_this_dag_will_not')

my_function()

파일에 접근하면 두 DAG 생성자 모두 호출되지만, dag_1 만 최상위 레벨(globals())에 존재하므로 이 DAG만 Airflow에 추가돼요. dag_2 는 로드되지 않습니다.

Note

DAG 번들 안에서 DAG를 검색할 때, Airflow는 최적화를 위해 airflowdag 문자열(대소문자 구분 없이)을 포함하는 Python 파일만 고려해요.

대신 모든 Python 파일을 고려하려면 DAG_DISCOVERY_SAFE_MODE 설정 플래그를 비활성화하면 돼요.

DAG 번들 안이나 그 하위 폴더에 .airflowignore 파일을 두면 로더가 무시할 파일 패턴을 지정할 수도 있어요. 이 파일은 자신이 있는 디렉터리와 그 아래 모든 하위 폴더를 대상으로 해요. 파일 문법의 자세한 내용은 아래 .airflowignore 섹션을 보면 돼요.

.airflowignore 가 요구사항을 충족하지 못하고, Python 파일을 Airflow가 파싱할지 여부를 더 유연하게 제어하고 싶다면 config 파일에 might_contain_dag_callable 을 설정해 원하는 콜러블(callable)을 꽂을 수 있어요. 이 콜러블은 기본 Airflow 휴리스틱(즉, Python 파일에 airflowdag 문자열이 대소문자 구분 없이 존재하는지 확인하는 것)을 대체한다는 점에 유의하세요.

def might_contain_dag(file_path: str, zip_file: zipfile.ZipFile | None = None) -> bool:
    # file_path 에 DAG가 정의되어 있는지 확인하는 로직
    # file_path 를 파싱해야 하면 True, 아니면 False 반환

DAG 실행 (Running Dags)

DAG는 두 가지 방식 중 하나로 실행돼요.

  • 수동으로 또는 API를 통해 trigger(트리거) 될 때
  • DAG의 일부로 정의된 schedule(스케줄) 에 따라

DAG가 스케줄을 반드시 가져야 하는 것은 아니지만, 정의하는 것이 아주 흔해요. schedule 인자로 정의하는데, 이렇게요.

with DAG("my_daily_dag", schedule="@daily"):
    ...

schedule 인자에 쓸 수 있는 값들은 다양해요.

with DAG("my_daily_dag", schedule="0 0 * * *"):
    ...

with DAG("my_one_time_dag", schedule="@once"):
    ...

with DAG("my_continuous_dag", schedule="@continuous"):
    ...

Tip

다양한 스케줄링 방식에 대한 더 자세한 내용은 Authoring and Scheduling 을 참고하세요.

DAG를 실행할 때마다 Airflow가 Dag Run이라고 부르는 새로운 인스턴스가 만들어져요. 같은 DAG에 대해 Dag Run들은 병렬로 실행될 수 있고, 각각은 data interval(데이터 구간) 을 가지는데, 이는 태스크가 다뤄야 할 데이터의 기간을 식별해 줘요.

이것이 왜 유용한지 예를 들어볼게요. 매일의 실험 데이터를 처리하는 DAG를 작성한다고 해보죠. 이 DAG를 재작성했는데, 지난 3개월간의 데이터에 대해 실행하고 싶다고요. 문제없어요. Airflow는 DAG를 backfill(역추적 실행) 해서 지난 3개월의 매일마다 DAG의 복사본을 한 번에 실행할 수 있으니까요.

이 Dag Run들은 모두 같은 실제 날짜에 시작됐지만, 각 Dag Run은 그 3개월 기간 중 하루를 덮는 data interval을 하나씩 가지게 돼요. 그리고 DAG 안의 모든 태스크·operator·sensor가 실행될 때 보는 것이 바로 그 data interval이에요.

DAG가 실행될 때마다 Dag Run으로 인스턴스화되는 것과 마찬가지로, DAG 안에 지정된 태스크도 함께 Task Instances로 인스턴스화돼요.

Dag Run은 시작할 때의 start date와 끝날 때의 end date를 가져요. 이 기간은 DAG가 실제로 '실행된' 시간을 나타내요. Dag Run의 start/end date 외에도 logical date(논리적 날짜) (이전에는 execution date라고 불림)라는 또 다른 날짜가 있는데, 이는 Dag Run이 스케줄되거나 트리거되기로 의도된 시간을 나타내요. 이것이 logical 이라고 불리는 이유는, Dag Run의 맥락에 따라 여러 의미를 갖는 추상적인 성격 때문이에요.

예를 들어, 사용자가 Dag Run을 수동으로 트리거했다면 그 logical date는 Dag Run이 트리거된 날짜·시간이 되고, 그 값은 Dag Run의 start date와 같아야 해요. 반면 DAG가 특정 schedule interval과 함께 자동으로 스케줄될 때는, logical date가 data interval의 시작을 표시하는 시간을 나타내요. 그러면 Dag Run의 start date는 logical date + scheduled interval 이 돼요.

DAG 할당 (Dag Assignment)

모든 Operator/Task는 실행되려면 반드시 DAG에 할당되어야 해요. Airflow가 DAG를 명시적으로 넘기지 않아도 계산하는 방법이 여러 가지 있어요.

  • with DAG 블록 안에서 Operator를 선언하는 경우
  • @dag 데코레이터 안에서 Operator를 선언하는 경우
  • DAG를 가진 Operator의 upstream 또는 downstream 위치에 Operator를 둔 경우

그 외에는 각 Operator에 dag= 로 넘겨줘야 해요.

기본 인자 (Default Arguments)

종종 DAG 안의 많은 Operator가 같은 기본 인자(예: retries)가 필요해요. 이걸 매 Operator마다 일일이 지정하는 대신, DAG를 만들 때 default_args 를 넘기면 그 DAG에 연결된 모든 operator에 자동으로 적용돼요.

import pendulum

with DAG(
    dag_id="my_dag",
    start_date=pendulum.datetime(2016, 1, 1),
    schedule="@daily",
    default_args={"retries": 2},
):
    op = BashOperator(task_id="hello_world", bash_command="Hello World!")
    print(op.retries)  # 2

DAG 데코레이터 (The Dag decorator)

버전 2.0에서 추가됨.

컨텍스트 매니저나 Dag() 생성자로 단일 DAG를 선언하는 전통적인 방식과 더불어, 함수에 @dag 를 붙이면 그 함수를 DAG 생성기 함수로 바꿀 수도 있어요.

from typing import TYPE_CHECKING, Any

import httpx
import pendulum

from airflow.providers.standard.operators.bash import BashOperator
from airflow.sdk import BaseOperator, dag, task

if TYPE_CHECKING:
    from airflow.sdk import Context


class GetRequestOperator(BaseOperator):
    """Custom operator to send GET request to provided url"""

    template_fields = ("url",)

    def __init__(self, *, url: str, **kwargs):
        super().__init__(**kwargs)
        self.url = url

    def execute(self, context: Context):
        return httpx.get(self.url).json()


@dag(
    schedule=None,
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    catchup=False,
    tags=["example"],
)
def example_dag_decorator(url: str = "https://httpbingo.org/get"):
    """
    DAG to get IP address and echo it via BashOperator.

    :param url: URL to get IP address from. Defaults to "https://httpbingo.org/get".
    """
    get_ip = GetRequestOperator(task_id="get_ip", url=url)

    @task(multiple_outputs=True)
    def prepare_command(raw_json: dict[str, Any]) -> dict[str, str]:
        external_ip = raw_json["origin"]
        try:
            ipaddress.ip_address(external_ip)
            return {
                "command": f"echo 'Seems like today your server executing Airflow is connected from IP {external_ip}'",
            }
        except ValueError:
            raise ValueError(f"Invalid IP address: '{external_ip}'.")

    command_info = prepare_command(get_ip.output)

    BashOperator(task_id="echo_ip_info", bash_command=command_info["command"])


example_dag = example_dag_decorator()

데코레이터는 DAG를 깔끔하게 만드는 새로운 방법일 뿐 아니라, 함수에 있는 파라미터들을 모두 DAG 파라미터로 설정해 줘서, DAG를 트리거할 때 이 파라미터들을 설정할 수 있게 해줘요. 그러면 Python 코드에서, 또는 Jinja 템플릿 안의 {{context.params}} 에서 파라미터에 접근할 수 있어요.

Note

Airflow는 DAG 파일의 최상위 레벨에 나타나는 DAG만 로드해요. 즉 @dag 로 함수만 선언해서는 안 되고, 위 예시처럼 DAG 파일에서 그 함수를 한 번 이상 호출해 최상위 객체에 할당해야 해요.

제어 흐름 (Control Flow)

기본적으로 DAG는, 의존하고 있는 모든 태스크가 성공했을 때만 태스크를 실행해요. 하지만 이를 수정하는 방법이 여러 가지 있어요.

  • Branching(분기) - 조건에 따라 어떤 태스크로 진행할지 선택하기
  • Trigger Rules(트리거 규칙) - DAG가 태스크를 실행할 조건을 설정하기
  • Setup and Teardown - setup/teardown 관계 정의하기
  • Latest Only(최신만) - 현재에 대해서 실행하는 DAG에서만 동작하는 특수한 분기 형태
  • Depends On Past(이전 실행 의존) - 태스크가 이전 실행의 자기 자신에 의존하게 하기

분기 (Branching)

branching을 사용하면 DAG가 의존 태스크를 전부 실행하지 않고, 하나 이상의 경로를 골라서 내려가도록 만들 수 있어요. 이때 @task.branch 데코레이터가 등장하지요.

@task.branch 데코레이터는 @task 와 매우 비슷한데, 데코레이트된 함수가 태스크의 ID(또는 ID 리스트)를 반환하리라고 기대한다는 점만 달라요. 지정된 태스크가 따라가게 되고, 나머지 모든 경로는 건너뛰어요(skip). 또한 None 을 반환해 하위의 모든 태스크를 건너뛸 수도 있어요.

Python 함수가 반환하는 task_id 는 @task.branch 로 데코레이트된 태스크의 직접 하위(downstream) 태스크를 참조해야 해요.

Note

어떤 태스크가 분기 operator의 하위이면서 동시에 선택된 태스크 하나 이상의 하위라면, 그 태스크는 건너뛰지 않아요.

../_images/branch_note.png

분기 태스크의 경로는 branch_a, join, branch_b 이에요. joinbranch_a 의 하위 태스크이므로, 분기 결정의 일부로 반환되지 않았어도 여전히 실행돼요.

@task.branch 는 XCom과 함께 사용할 수도 있어서, 분기 컨텍스트가 upstream 태스크에 기반해 어떤 분기를 따라갈지 동적으로 결정할 수 있게 해줘요. 예를 들어보면요.

@task.branch(task_id="branch_task")
def branch_func(ti=None):
    xcom_value = int(ti.xcom_pull(task_ids="start_task"))
    if xcom_value >= 5:
        return "continue_task"
    elif xcom_value >= 3:
        return "stop_task"
    else:
        return None


start_op = BashOperator(
    task_id="start_task",
    bash_command="echo 5",
    do_xcom_push=True,
    dag=dag,
)

branch_op = branch_func()

continue_op = EmptyOperator(task_id="continue_task", dag=dag)
stop_op = EmptyOperator(task_id="stop_task", dag=dag)

start_op >> branch_op >> [continue_op, stop_op]

분기 기능을 갖춘 자신만의 operator를 구현하고 싶다면 BaseBranchOperator 를 상속하면 돼요. 이 클래스는 @task.branch 데코레이터와 비슷하게 동작하지만, choose_branch 메서드의 구현을 직접 제공해야 해요.

Note

DAG에서 BranchPythonOperator 를 직접 인스턴스화하는 것보다 @task.branch 데코레이터를 권장해요. 전자는 일반적으로 커스텀 operator를 구현하려고 서브클래싱할 때만 사용해야 해요.

@task.branch 의 콜러블과 마찬가지로 이 메서드는 하위 태스크의 ID나 태스크 ID 리스트를 반환할 수 있어요. 반환된 것들이 실행되고 나머지는 모두 건너뛰게 되지요. 또한 None 을 반환해 하위의 모든 태스크를 건너뛸 수도 있어요.

class MyBranchOperator(BaseBranchOperator):
    def choose_branch(self, context):
        """
        Run an extra branch on the first day of the month
        """
        if context['data_interval_start'].day == 1:
            return ['daily_task_id', 'monthly_task_id']
        elif context['data_interval_start'].day == 2:
            return 'daily_task_id'
        else:
            return None

일반 Python 코드용 @task.branch 데코레이터와 비슷하게, 가상 환경을 사용하는 @task.branch_virtualenv 나 외부 Python을 사용하는 @task.branch_external_python 이라는 분기 데코레이터도 있어요.

최신만 실행 (Latest Only)

Airflow의 Dag Run은 종종 현재 날짜와 다른 날짜를 위해 실행돼요. 예를 들어 데이터를 역추적(backfill)하느라 지난달의 매일마다 DAG 복사본을 하나씩 실행하는 경우가 있지요.

하지만 이전 날짜를 위해 DAG의 일부(또는 전체)를 실행하고 싶지 않은 상황도 있어요. 이때 LatestOnlyOperator 를 사용할 수 있어요.

이 특수 Operator는, 가장 최신 Dag Run이 아니라면(지금의 벽시계 시간이 그 Dag run의 execution_time과 다음 scheduled execution_time 사이에 있고, 외부 트리거된 run이 아닌 경우) 자기 하위의 모든 태스크를 건너뜁니다.

예시를 볼게요.

import datetime

import pendulum

from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.providers.standard.operators.latest_only import LatestOnlyOperator
from airflow.sdk import DAG, TriggerRule

with DAG(
    dag_id="latest_only_with_trigger",
    schedule=datetime.timedelta(hours=4),
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    catchup=False,
    tags=["example", "example3"],
) as dag:
    latest_only = LatestOnlyOperator(task_id="latest_only")
    task1 = EmptyOperator(task_id="task1")
    task2 = EmptyOperator(task_id="task2")
    task3 = EmptyOperator(task_id="task3")
    task4 = EmptyOperator(task_id="task4", trigger_rule=TriggerRule.ALL_DONE)

    latest_only >> task1 >> [task3, task4]
    task2 >> [task3, task4]

이 DAG의 경우를 살펴보면요.

  • task1latest_only 의 직접 하위라서, 가장 최신 run을 제외한 모든 run에서 건너뛰어져요.
  • task2latest_only 와 완전히 독립적이라 모든 스케줄된 기간에 실행돼요.
  • task3task1task2 의 하위인데, 기본 트리거 규칙이 all_success 이므로 task1 의 건너뜀(skip)이 연쇄적으로 전파됩니다.
  • task4task1task2 의 하위이지만, trigger_ruleall_done 으로 설정되어 있어서 건너뛰지 않아요.

../_images/latest_only_with_trigger.png

이전 실행 의존 (Depends On Past)

태스크가 이전 Dag Run에서 이전 실행이 성공했을 때만 실행되게 할 수도 있어요. 이를 사용하려면 태스크의 depends_on_past 인자를 True 로 설정하기만 하면 돼요.

DAG가 생애의 아주 처음, 특히 첫 자동화된 run일 때는 의존할 이전 run이 없으므로 태스크가 여전히 실행된다는 점에 유의하세요.

트리거 규칙 (Trigger Rules)

기본적으로 Airflow는 태스크를 실행하기 전에 그 태스크의 모든 upstream(직접 상위) 태스크가 성공(successful)하기를 기다려요.

하지만 이것은 어디까지나 기본 동작일 뿐이고, 태스크의 trigger_rule 인자로 제어할 수 있어요. trigger_rule 의 옵션은 다음과 같아요.

  • all_success (기본값): 모든 upstream 태스크가 성공함
  • all_failed: 모든 upstream 태스크가 failed 또는 upstream_failed 상태임
  • all_done: 모든 upstream 태스크가 실행을 마침
  • all_done_setup_success: all_done 과 같지만, 태스크에 upstream setup 태스크가 있으면 그중 하나 이상이 성공해야 함. teardown 태스크의 기본 트리거 규칙임
  • all_done_min_one_success: 건너뛰지 않은 모든 upstream 태스크가 실행을 마치고, upstream 태스크 중 하나 이상이 성공함
  • all_skipped: 모든 upstream 태스크가 skipped 상태임
  • one_failed: upstream 태스크 중 하나 이상이 실패함 (모든 upstream 태스크가 끝나기를 기다리지 않음)
  • one_success: upstream 태스크 중 하나 이상이 성공함 (모든 upstream 태스크가 끝나기를 기다리지 않음)
  • one_done: upstream 태스크 중 하나 이상이 성공하거나 실패함
  • none_failed: 모든 upstream 태스크가 failedupstream_failed 상태가 아님, 즉 모두 성공했거나 건너뛰어졌음
  • none_failed_min_one_success: 모든 upstream 태스크가 failedupstream_failed 상태가 아니고, upstream 태스크 중 하나 이상이 성공함
  • none_skipped: upstream 태스크 중 어떤 것도 skipped 상태가 아님, 즉 모두 success, failed, upstream_failed, 또는 removed 상태임
  • always: 의존성이 전혀 없음, 언제든 이 태스크 실행

Note

Task Instance 상태 removed 는 run이 시작된 이후 태스크가 DAG에서 사라졌다는 뜻이에요. 트리거 규칙에서 removed 는 종결(terminal) 상태로 취급돼요. all_done, all_done_setup_success, all_done_min_one_success 같은 규칙에서 "done"으로 세어지지만, success, failed, upstream_failed, skipped 로는 세어지지 않아요. 동적으로 매핑된(mapped) 태스크의 경우, removed upstream 맵 인덱스는 all_success, all_failed, none_failed, none_failed_min_one_success, all_done_min_one_success 의 실패 카운트에서도 차감돼요.

이 트리거 규칙은 원한다면 Depends On Past 기능과도 결합할 수 있어요.

Note

트리거 규칙과 건너뛴 태스크(skipped) 사이의 상호작용, 특히 분기(branching) 작업의 일부로 건너뛰어진 태스크를 잘 인지하는 것이 중요해요. 분기 작업의 하위에서 all_success 나 all_failed 를 쓰는 것은 거의 항상 피해야 합니다.

건너뛴 태스크는 all_successall_failed 트리거 규칙을 통해 연쇄적으로 전파되어, 그 태스크들도 함께 건너뛰게 만들어요. 다음과 같은 DAG를 생각해 보죠.

# dags/branch_without_trigger.py
import pendulum

from airflow.sdk import task
from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator

dag = DAG(
    dag_id="branch_without_trigger",
    schedule="@once",
    start_date=pendulum.datetime(2019, 2, 28, tz="UTC"),
)

run_this_first = EmptyOperator(task_id="run_this_first", dag=dag)


@task.branch(task_id="branching")
def do_branching():
    return "branch_a"


branching = do_branching()

branch_a = EmptyOperator(task_id="branch_a", dag=dag)
follow_branch_a = EmptyOperator(task_id="follow_branch_a", dag=dag)

branch_false = EmptyOperator(task_id="branch_false", dag=dag)

join = EmptyOperator(task_id="join", dag=dag)

run_this_first >> branching
branching >> branch_a >> follow_branch_a >> join
branching >> branch_false >> join

joinfollow_branch_abranch_false 의 하위예요. join 태스크는 trigger_rule 이 기본값으로 all_success 로 설정되어 있어서, 분기 작업으로 생긴 skip이 all_success 로 표시된 태스크를 건너뛰게끔 연쇄 전파되므로 skipped로 나타나요.

../_images/branch_without_trigger.png

join 태스크에서 trigger_rulenone_failed_min_one_success 로 설정하면, 의도한 동작을 얻을 수 있어요.

../_images/branch_with_trigger.png

Setup과 teardown

데이터 워크플로우에서는 리소스(예: 컴퓨팅 리소스)를 만들고, 그것을 사용해 일을 한 다음, 다시 내려놓는(tear down) 것이 흔해요. Airflow는 이런 필요를 지원하는 setup 및 teardown 태스크를 제공해요.

이 기능을 사용하는 방법에 대한 자세한 내용은 본문 기사 Setup and Teardown 을 참고하세요.

동적 DAG (Dynamic Dags)

DAG는 Python 코드로 정의되기 때문에 순수 선언적일 필요가 없어요. 루프, 함수 등을 자유롭게 사용해 DAG를 정의할 수 있지요.

예를 들어, for 루프를 사용해 태스크를 정의하는 DAG가 있어요.

with DAG("loop_example", ...):
    first = EmptyOperator(task_id="first")
    last = EmptyOperator(task_id="last")

    options = ["branch_a", "branch_b", "branch_c", "branch_d"]
    for option in options:
        t = EmptyOperator(task_id=option)
        first >> t >> last

일반적으로 DAG 태스크의 토폴로지(배치 구조)는 상대적으로 안정적으로 유지하라고 권장해요. 동적 DAG는 보통 설정 옵션을 동적으로 로드하거나 operator 옵션을 변경하는 데 더 적합해요.

DAG 시각화 (Dag Visualization)

DAG를 시각적으로 표현하고 싶다면 두 가지 방법이 있어요.

  • Airflow UI를 열고 DAG로 이동한 뒤 "Graph"를 선택하기
  • airflow dags show 를 실행해 이미지 파일로 렌더링하기

일반적으로 Graph 뷰를 권장해요. 선택한 Dag Run 안의 모든 Task Instances의 상태도 함께 보여주기 때문이에요.

물론 DAG를 개발하면서 점점 복잡해지기 마련이에요. 그래서 이런 DAG 뷰를 더 이해하기 쉽게 수정할 수 있는 몇 가지 방법을 제공해요.

TaskGroups

TaskGroup은 Graph 뷰에서 태스크들을 계층적 그룹으로 정리하는 데 사용돼요. 반복되는 패턴을 만들고 시각적 복잡도를 줄이는 데 유용하지요.

TaskGroup 안의 태스크들은 원래의 같은 DAG에 존재하며, 모든 DAG 설정과 pool 구성을 그대로 따릅니다.

See also

TaskGrouptask_group 의 API 레퍼런스

../_images/task_group.gif

>><< 연산자로 TaskGroup 안의 모든 태스크에 걸쳐 의존 관계를 적용할 수 있어요. 예를 들어, 아래 코드는 task1task2 를 TaskGroup group1 에 넣고 두 태스크 모두를 task3 의 upstream으로 둬요.

from airflow.sdk import task_group


@task_group()
def group1():
    task1 = EmptyOperator(task_id="task1")
    task2 = EmptyOperator(task_id="task2")


task3 = EmptyOperator(task_id="task3")

group1() >> task3

TaskGroup도 DAG처럼 default_args 를 지원하는데, DAG 레벨의 default_args 를 덮어씁니다.

import datetime

from airflow.sdk import DAG
from airflow.sdk import task_group
from airflow.providers.standard.operators.bash import BashOperator
from airflow.providers.standard.operators.empty import EmptyOperator

with DAG(
    dag_id="dag1",
    start_date=datetime.datetime(2016, 1, 1),
    schedule="@daily",
    default_args={"retries": 1},
):

    @task_group(default_args={"retries": 3})
    def group1():
        """This docstring will become the tooltip for the TaskGroup."""
        task1 = EmptyOperator(task_id="task1")
        task2 = BashOperator(task_id="task2", bash_command="echo Hello World!", retries=2)
        print(task1.retries)  # 3
        print(task2.retries)  # 2

TaskGroup의 더 고급 사용법을 보고 싶다면 Airflow에 포함된 example_task_group_decorator.py 예제 DAG를 살펴볼 수 있어요.

Note

기본적으로 자식 태스크/하위 TaskGroup의 ID는 상위 TaskGroup의 group_id로 접두사(prefix)가 붙어요. 이렇게 해서 DAG 전체에서 group_id와 task_id의 고유성을 보장하지요.

접두사를 비활성화하려면 TaskGroup을 만들 때 prefix_group_id=False 를 넘기면 되는데, 이제 모든 태스크와 그룹에 각자 고유한 ID가 있는지 직접 보장해야 한다는 점을 유의하세요.

Note

@task_group 데코레이터를 사용할 때, 데코레이트된 함수의 docstring은 tooltip 값이 명시적으로 제공되지 않는 한 UI에서 TaskGroup의 툴팁으로 사용돼요.

엣지 라벨 (Edge Labels)

태스크를 그룹으로 묶는 것 외에도, Graph 뷰에서 서로 다른 태스크 간의 의존성 엣지(dependency edge) 에 라벨을 붙일 수 있어요. 이것은 DAG의 분기 부분에서 특히 유용한데, 특정 분기가 언제 실행될 수 있는지의 조건을 라벨로 표시할 수 있으니까요.

라벨을 추가하려면 >><< 연산자에 인라인으로 직접 쓰면 돼요.

from airflow.sdk import Label

my_task >> Label("When empty") >> other_task

아니면 Label 객체를 set_upstream/set_downstream 에 넘길 수도 있어요.

from airflow.sdk import Label

my_task.set_downstream(other_task, Label("When empty"))

여러 분기에 라벨을 붙이는 것을 보여주는 예제 DAG가 있어요.

../_images/edge_label_example.png

with DAG(
    "example_branch_labels",
    schedule="@daily",
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    catchup=False,
    tags=["example"],
) as dag:
    ingest = EmptyOperator(task_id="ingest")
    analyse = EmptyOperator(task_id="analyze")
    check = EmptyOperator(task_id="check_integrity")
    describe = EmptyOperator(task_id="describe_integrity")
    error = EmptyOperator(task_id="email_error")
    save = EmptyOperator(task_id="save")
    report = EmptyOperator(task_id="report")

    ingest >> analyse >> check
    check >> Label("No errors") >> save >> report
    check >> Label("Errors found") >> describe >> error >> report

DAG, Task Group & Task 문서화 (Dag, Task Group & Task Documentation)

DAG, TaskGroup, 태스크 객체에 웹 인터페이스에서 보이는 문서나 메모를 추가할 수 있어요.

정의하면 풍부한 콘텐츠(rich content)로 렌더링되는 특수 태스크 속성들이 몇 가지 있어요.

attribute rendered to
doc monospace
doc_json json
doc_yaml yaml
doc_md markdown
doc_rst reStructuredText

DAG와 TaskGroup의 경우 doc_md 만 해석되는 속성이라는 점에 유의하세요. doc_md 는 문자열이나 markdown 파일의 참조를 담을 수 있어요. markdown 파일은 .md 로 끝나는 문자열로 인식돼요. 상대 경로가 주어지면 Airflow Scheduler나 Dag 파서가 시작된 경로를 기준으로 로드돼요. markdown 파일이 존재하지 않으면 전달된 파일 이름이 텍스트로 사용되고, 예외는 표시되지 않아요. markdown 파일은 Dag 파싱 중에 로드되므로, markdown 내용의 변경 사항이 표시되려면 Dag 파싱 사이클 한 번이 걸립니다.

이 기능은 특히 태스크가 설정 파일에서 동적으로 만들 때 유용해요. Airflow에서 관련 태스크로 이끈 설정을 노출할 수 있으니까요.

"""
### My great Dag
"""

import pendulum
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.sdk import DAG, TaskGroup

with DAG(
    "my_dag",
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    schedule="@daily",
    catchup=False,
) as dag:
    dag.doc_md = __doc__

    t = EmptyOperator(task_id="foo")
    t.doc_md = """\
    #Title"
    Here's a [url](www.airbnb.com)
    """

    with TaskGroup("extract", doc_md="### Extract tasks"):
        EmptyOperator(task_id="extract_orders")

DAG 패키징 (Packaging Dags)

간단한 DAG는 보통 단일 Python 파일에 담기지만, 더 복잡한 DAG는 여러 파일에 걸쳐 있고 함께 배포되어야 할("vendored") 의존성을 가진 경우가 드물지 않아요.

이 모든 것을 표준 파일시스템 구조를 가진 DAG 번들 안에서 처리할 수도 있고, DAG와 모든 Python 파일을 단일 zip 파일로 패키징할 수도 있어요. 예를 들어, 두 DAG와 그들이 필요로 하는 의존성을 zip 파일로 배포한다고 하면 내용물은 다음과 같아요.

my_dag1.py
my_dag2.py
package1/__init__.py
package1/functions.py

패키징된 DAG에는 몇 가지 주의사항(caveat)이 붙어요.

  • 직렬화를 위해 pickling을 활성화한 경우 사용할 수 없어요
  • 컴파일된 라이브러리(예: libz.so)를 담을 수 없고, 순수 Python만 가능해요
  • Python의 sys.path 에 삽입되어 Airflow 프로세스의 다른 코드에서 import할 수 있으므로, 패키지 이름이 시스템에 이미 설치된 다른 패키지와 충돌하지 않도록 해야 해요

일반적으로 복잡한 컴파일된 의존성과 모듈 세트가 있다면 Python virtualenv 시스템을 사용하고 필요한 패키지를 pip 로 대상 시스템에 설치하는 것이 더 나을 가능성이 높아요.

.airflowignore

.airflowignore 파일은 DAG 번들이나 PLUGINS_FOLDER에서 Airflow가 의도적으로 무시해야 할 디렉터리 또는 파일을 지정해요. Airflow는 DAG_IGNORE_FILE_SYNTAX 설정 파라미터(Airflow 2.3에서 추가됨)로 지정되는 두 가지 문법 방식을 지원해요. regexpglob 이에요.

Note

기본 DAG_IGNORE_FILE_SYNTAX 는 Airflow 3 이상에서 glob 이에요(이전 버전에서는 regexp 였죠).

glob 문법(기본값)에서는 패턴이 .gitignore 파일의 것과 똑같이 동작해요.

  • * 문자는 / 를 제외한 임의 개수의 문자와 매칭돼요
  • ? 문자는 / 를 제외한 임의의 단일 문자와 매칭돼요
  • 범위 표기법, 예: [a-zA-Z] 는 범위 안의 문자 중 하나와 매칭하는 데 쓸 수 있어요
  • ! 접두사를 붙이면 패턴을 부정(negate)할 수 있어요. 패턴은 순서대로 평가되므로, 부정이 같은 파일에서 이전에 정의된 패턴이나 상위 디렉터리에 정의된 패턴을 덮어쓸 수 있어요
  • 이중 별표(**)는 디렉터리를 가로질러 매칭하는 데 쓸 수 있어요. 예: **/__pycache__/ 는 무한한 깊이의 각 하위 디렉터리의 __pycache__ 디렉터리를 무시해요
  • 패턴의 시작이나 중간(또는 둘 다)에 / 가 있으면 그 패턴은 해당 .airflowignore 파일 자신의 디렉터리 레벨에 상대적이에요. 그렇지 않으면 패턴은 .airflowignore 레벨 아래의 어떤 레벨에서도 매칭될 수 있어요

regexp 패턴 문법에서는 .airflowignore 의 각 줄이 정규식 패턴을 지정하고, 이름(Dag id가 아니라)이 패턴 중 하나라도 매칭되는 디렉터리나 파일은 무시돼요(내부적으로 Pattern.search() 로 매칭합니다). # 문자는 주석을 나타내는 데 쓰이며, # 로 시작하는 줄의 모든 문자는 무시돼요.

.airflowignore 파일은 DAG 번들에 넣어야 해요. 예를 들어 glob 문법으로 .airflowignore 파일을 준비할 수 있어요.

**/*project_a*
tenant_[0-9]*

그러면 DAG 번들 안의 project_a_dag_1.py, TESTING_project_a.py, tenant_1.py, project_a/dag_1.py, tenant_1/dag_1.py 같은 파일들이 무시돼요(디렉터리 이름이 패턴 중 하나와 매칭되면 그 디렉터리와 모든 하위 폴더는 Airflow가 전혀 스캔하지 않아요. 이는 DAG 찾기의 효율을 높여 줍니다).

.airflowignore 의 범위는 자신이 있는 디렉터리와 그 모든 하위 폴더예요. DAG 번들의 하위 폴더용으로도 .airflowignore 파일을 준비할 수 있으며, 그 경우 해당 하위 폴더에만 적용돼요.

DAG 의존성 (Dag Dependencies)

Airflow 2.1에서 추가됨

DAG 안의 태스크 사이의 의존성은 upstream·downstream 관계로 명시적으로 정의되는 반면, DAG 사이의 의존성은 조금 더 복잡해요. 일반적으로 한 DAG가 다른 DAG에 의존하는 방식은 두 가지가 있어요.

추가적인 어려움은, 한 DAG가 서로 다른 data interval을 가진 다른 DAG의 여러 실행을 기다리거나 트리거할 수 있다는 점이에요. 이런 의존성은 Dag 직렬화(serialization) 동안 스케줄러가 계산해요.

의존성 감지기(dependency detector)는 설정 가능해서, DependencyDetector 의 기본값과 다른 자신만의 로직을 구현할 수 있어요.

DAG 일시정지, 비활성화, 삭제 (Dag pausing, deactivation and deletion)

DAG는 "실행되지 않는" 상태에 관해 여러 상태를 가져요. DAG는 일시정지(paused)될 수 있고, 비활성화(deactivated)될 수 있으며, 마지막으로 DAG의 모든 메타데이터가 삭제될 수 있어요.

DAG는 DAGS_FOLDER 에 존재하고 스케줄러가 데이터베이스에 저장했지만, 사용자가 UI를 통해 비활성화를 선택한 경우 UI를 통해 일시정지될 수 있어요. "pause" 와 "unpause" 작업은 UI와 API로 가능해요. 일시정지된 DAG는 Scheduler가 스케줄링하지 않지만, UI에서 수동 실행을 위해 트리거할 수는 있어요. UI에서 일시정지된 DAG를 볼 수 있는데(Paused 탭) 일시정지되지 않은 DAG는 Active 탭에서 찾을 수 있어요. DAG가 일시정지되면 실행 중인 태스크는 완료가 허용되고, 모든 하위 태스크는 "Scheduled" 상태로 배치돼요. DAG가 일시정지 해제되면 "scheduled" 태스크들이 DAG 로직에 따라 실행되기 시작해요. "scheduled" 태스크가 없는 DAG는 자신의 스케줄에 따라 실행되기 시작합니다.

DAG는 DAGS_FOLDER 에서 제거함으로써 비활성화될 수 있어요(UI의 Active 탭과 혼동하지 마세요). 스케줄러가 DAGS_FOLDER 를 파싱할 때 이전에 보고 데이터베이스에 저장했던 DAG가 없다면, 그 DAG를 비활성화된 것으로 설정해요. 비활성화된 DAG의 메타데이터와 이력은 보존되고, DAG가 DAGS_FOLDER 에 다시 추가되면 다시 활성화되어 이력이 보이게 돼요. UI나 API로는 DAG를 활성화/비활성화할 수 없고, 오직 DAGS_FOLDER 에서 파일을 제거하는 것으로만 가능해요. 다시 말하지만, 스케줄러가 비활성화해도 DAG의 과거 실행 데이터는 전혀 손실되지 않아요. UI의 Active 탭은 Activated 이면서 Not paused 인 DAG를 가리킨다는 점에 유의하세요. 그래서 처음에는 약간 혼동될 수 있어요.

UI에서는 비활성화된 DAG를 볼 수 없어요. 가끔 과거 실행은 볼 수 있지만, 그 정보를 보려고 하면 DAG가 없다는 오류가 표시돼요.

UI나 API를 사용해 메타데이터 데이터베이스에서 DAG 메타데이터를 삭제할 수도 있는데, 이것이 항상 UI에서 DAG가 사라지는 결과를 만들지는 않아요. 이것도 처음에는 혼동될 수 있지요. DAG가 여전히 DAGS_FOLDER 에 있을 때 메타데이터를 삭제하면, 스케줄러가 폴더를 파싱하면서 DAG가 다시 나타나게 되고, 단지 DAG의 과거 실행 정보만 제거돼요.

즉, DAG와 모든 과거 메타데이터를 실제로 삭제하려면 세 단계를 거쳐야 해요.

  • DAG를 일시정지(pause)한다
  • UI나 API를 통해 데이터베이스에서 과거 메타데이터를 삭제한다
  • DAGS_FOLDER 에서 DAG 파일을 삭제하고 비활성화될 때까지 기다린다

DAG 자동 일시정지 (Dag Auto-pausing) (실험적)

DAG는 자동 일시정지되도록 설정할 수도 있어요. DAG가 N 번 연속 실패하면 자동으로 비활성화할 수 있게 해주는 Airflow 설정이 있어요.

DAG 인자로도 이 설정을 제공하고 덮어쓸 수 있어요.

데드라인 알림 (Deadline Alerts)

버전 3.1에서 추가됨.

Deadline Alerts를 사용하면 Dag run에 시간 임계값을 설정하고, 그 임계값을 초과할 때 자동으로 대응할 수 있어요. 고정 datetime에 상대적으로 데드라인을 설정하거나, 사용 가능한 계산된 참조(예: Dag 대기 시간이나 시작 시간) 중 하나를 사용하거나, 자신만의 커스텀 참조를 구현할 수도 있어요. 데드라인이 초과되면 알려주거나 다른 조치를 취할 수 있는 콜백이 트리거돼요.

기존 이메일 Notifier를 사용하는 간단한 예시가 있어요.

from datetime import timedelta
from airflow import DAG
from airflow.providers.smtp.notifications.smtp import SmtpNotifier
from airflow.sdk.definitions.deadline import DeadlineAlert, DeadlineReference

with DAG(
    dag_id="email_deadline",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_QUEUED_AT,
        interval=timedelta(minutes=30),
        callback=SmtpNotifier(
            to="[email protected]",
            subject="🚨 Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}",
            html_content="The Dag Run {{ dag_run.dag_run_id }} has been running for more than 30 minutes since being queued.",
        ),
    ),
):
    EmptyOperator(task_id="task1")

이 예시는 DAG가 대기열에 들어간 후 30분 안에 끝나지 않으면 이메일 알림을 보내요.

Deadline Alerts 구현과 구성에 대한 자세한 내용은 Deadline Alerts 를 참고하세요.

DAG 테스트 (Testing a Dag)

Dag 코드의 문법/기능을 검증할 수 있는 편리한 옵션이 있어요.

dag 객체에는 test() 함수가 있어서, 직접 또는 테스트 스위트 안에서 호출할 수 있어요.

시뮬레이션된 dag run (Simulated dag run)

dag 의 유효성을 확인하는 가장 간단한 방법은 dag 정의 모듈의 끝에 다음 스니펫을 쓰는 거예요.

if __name__ == "__main__":
    dag.test()

그리고 모듈을 Python 스크립트로 실행하면 돼요(AIRFLOW_HOME 같은 충분한 환경이 제공된다고 가정해요).

dag.test() 호출은 LocalExecutor 실행 세션과 유사한 시뮬레이션된 실행 흐름을 호출해서, Airflow 제출에서처럼 Dag 내용을 실행해요.

실제 실행 흐름 (Real execution flow)

실제 Airflow Executor에 대해 Dag를 테스트하고 싶다면 같은 메커니즘을 사용할 수 있어요. 호출에 use_executor 플래그를 지정하면, 현재 적용된 Airflow Configuration의 Airflow Executor가 호출되어 Dag의 워크로드를 실행해요.

dag.test(use_executor=True)

pytest 테스트 스위트 안에서 이 호출을 사용하면, 미리 정의된 설정 값을 쉽게 바꿔가며 쓸 수 있는 Airflow pytest 플러그인의 conf_vars 픽스처를 활용할 수 있어요.