Dags

Dags

워크플로우 실행에 필요한 모든 것을 담는 모델인 Dag를 설명하는 문서예요. Dag 선언 방법, Task 의존성, 로딩, 실행, 기본 인자, @dag 데코레이터, 제어 흐름(Branching·Trigger Rules·Latest Only·Depends On Past), 동적 Dag, 시각화, 패키징, .airflowignore, Dag 의존성, pause·비활성화·삭제, Deadline Alerts, 테스트까지 폭넓게 다뤄요.

출처: 문서

본문

Dag는 워크플로우 실행에 필요한 모든 것을 캡슐화하는 모델이에요. 몇 가지 Dag 속성으로는 다음이 있어요:

  • 스케줄 (Schedule): 워크플로우가 언제 실행되어야 하는지.
  • 태스크 (Tasks): tasks는 워커에서 실행되는 개별 작업 단위.
  • 태스크 의존성 (Task Dependencies): tasks가 실행되는 순서와 조건.
  • 콜백 (Callbacks): 전체 워크플로우가 완료될 때 취할 동작.
  • 추가 파라미터 (Additional Parameters): 그 외 여러 운영 세부 사항.

다음은 기본 예시 Dag예요:

A, B, C, D 네 개의 Task를 정의하고, 그것들이 실행되어야 하는 순서와 어떤 태스크가 무엇에 의존하는지를 지시해요. 또한 Dag를 얼마나 자주 실행할지도 말해요 — "내일부터 5분마다", 또는 "2020년 1월 1일부터 매일"처럼요.

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")

또는 표준 생성자를 사용해 사용하는 모든 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 생성기로 바꿀 수 있어요:

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()

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

Task 의존성 (Task Dependencies)

Task/Operator는 보통 혼자 존재하지 않아요. 다른 태스크(그 업스트림)에 의존하고, 다른 태스크가 그것에 의존해요(그 다운스트림). 태스크 간의 이러한 의존성을 선언하는 것이 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

# Replaces
# [op1, op2] >> op3
# [op1, op2] >> op4
cross_downstream([op1, op2], [op3, op4])

그리고 의존성을 사슬처럼 연결하고 싶다면 chain을 사용할 수 있어요:

from airflow.sdk import chain

# Replaces op1 >> op2 >> op3 >> op4
chain(op1, op2, op3, op4)

# You can also do it dynamically
chain(*[EmptyOperator(task_id=f"op{i}") for i in range(1, 6)])

Chain은 같은 크기의 목록에 대해 쌍별(pairwise) 의존성을 만들 수도 있어요(cross_downstream이 만드는 교차 의존성(cross dependencies) 과는 다르다는 점 참고!):

from airflow.sdk import chain

# Replaces
# op1 >> op2 >> op4 >> op6
# op1 >> op3 >> op5 >> op6
chain(op1, [op2, op3], [op4, op5], op6)

Dags 로딩 (Loading Dags)

Airflow는 Dag bundle의 Python 소스 파일에서 Dags를 로드해요. 각 파일을 가져와 실행한 다음, 그 파일에서 어떤 Dag 객체든 로드해요.

이것은 Python 파일당 여러 Dags를 정의할 수 있거나, import를 사용해 매우 복잡한 Dag 하나를 여러 Python 파일에 걸쳐 펼칠 수도 있다는 뜻이에요.

다만 Airflow가 Python 파일에서 Dags를 로드할 때는 최상위 레벨 에서 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())에 있으므로 Airflow에 추가되는 것은 그것뿐이에요. dag_2는 로드되지 않아요.

참고 (Note)

Dag bundle 안에서 Dags를 검색할 때 Airflow는 최적화로 airflowdag 문자열(대소문자 무시)을 포함하는 Python 파일만 고려해요.

모든 Python 파일을 고려하려면 [core] dag_discovery_safe_mode 구성 플래그를 False로 설정하세요. 이는 source에 그 문자열을 포함하지 않는 wrapper나 추상화를 통해 Dags를 정의할 때 사용하는 설정이에요. 이 플래그는 Dag 파일 processor가 읽으므로 해당 컴포넌트(및 airflow dags reserialize를 실행하는 것)에 설정하고 변경 사항을 적용하기 위해 재시작하세요 — 그렇지 않으면 수동 reserialize 후 Dags가 나타났다가 다음 processor 스캔에서 다시 사라질 수 있어요.

Dag bundle(또는 그 하위 폴더) 안에 .airflowignore 파일을 제공할 수도 있는데, 이는 로더가 무시할 파일들의 패턴을 설명해요. 그것은 있는 디렉토리와 그 아래의 모든 하위 폴더를 포함해요. 파일 문법에 대한 자세한 내용은 아래 .airflowignore를 참고하세요.

.airflowignore가 요구에 맞지 않고, Python 파일이 airflow에 의해 파싱되어야 할지를 제어하는 더 유연한 방법을 원한다면, 구성 파일에서 might_contain_dag_callable을 설정해 자신의 callable을 연결할 수 있어요. 참고로 이 callable은 기본 Airflow 휴리스틱, 즉 python 파일에 airflowdag 문자열(대소문자 무시)이 있는지 확인하는 것을 대체할 거예요.

def might_contain_dag(file_path: str, zip_file: zipfile.ZipFile | None = None) -> bool:
    # Your logic to check if there are Dags defined in the file_path
    # Return True if the file_path needs to be parsed, otherwise False

Dags 실행 (Running Dags)

Dags는 두 가지 방식 중 하나로 실행돼요:

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

Dags는 스케줄을 요구하지 않지만 정의하는 것이 매우 흔해요. 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 Runs는 같은 Dag에 대해 병렬로 실행될 수 있고, 각각은 태스크가 작업해야 할 데이터의 기간을 식별하는 정의된 데이터 구간(data interval)을 가져요.

이것이 왜 유용한지에 대한 예로, 매일의 실험 데이터 세트를 처리하는 Dag를 작성한다고 생각해 봐요. 재작성되었고 이전 3개월 데이터에 실행하고 싶다면 문제없어요. Airflow가 Dag를 backfill 하고 이전 3개월의 매일마다 사본을 한 번에 실행할 수 있기 때문이에요.

그 Dag Runs는 모두 같은 실제 날짜에 시작됐겠지만, 각 Dag Run은 그 3개월 기간의 하루를 커버하는 하나의 데이터 구간을 가지며, 그리고 Dag 내부의 모든 태스크·operator·sensor가 실행될 때 보는 것이 바로 그 데이터 구간이에요.

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

Dag run은 시작할 때 start date, 끝날 때 end date를 가져요. 이 기간은 Dag가 실제로 '실행'된 시간을 설명해요. Dag run의 start date와 end date 외에, logical date(이전에는 execution date로 알려짐)라는 다른 날짜가 있어요. 이는 Dag run이 스케줄되거나 트리거된 의도된 시간을 설명해요. logical이라고 불리는 이유는 Dag run 자체의 컨텍스트에 따라 여러 의미를 갖는 추상적인 특성 때문이에요.

예를 들어 Dag run이 사용자에 의해 수동 트리거되면, 그 logical date는 Dag run이 트리거된 날짜와 시간이며, 값은 Dag run의 start date와 같아야 해요. 하지만 Dag가 특정 스케줄 간격으로 자동 스케줄될 때, logical date는 데이터 구간의 시작을 표시하는 시간을 나타내고, Dag run의 start date는 그때 logical date + 스케줄 간격이 돼요.

팁 (Tip)

logical date에 대한 자세한 내용은 Data IntervalWhat does execution_date mean?을 참고하세요.

Dag 할당 (Dag Assignment)

모든 Operator/Task가 실행되려면 Dag에 할당되어야 한다는 점을 참고하세요. Airflow는 명시적으로 전달하지 않아도 Dag를 계산하는 몇 가지 방법이 있어요:

  • with DAG 블록 안에서 Operator를 선언하는 경우
  • @dag 데코레이터 안에서 Operator를 선언하는 경우
  • Dag가 있는 Operator의 업스트림이나 다운스트림에 Operator를 두는 경우

그 외에는 dag=로 각 Operator에 전달해야 해요.

기본 인자 (Default Arguments)

종종 Dag 안의 많은 Operators가 같은 기본 인자 세트(예: retries)를 필요로 해요. 매번 Operator마다 개별로 지정하는 대신, Dag를 만들 때 default_args를 전달하면 그것에 연결된 어떤 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 버전에 추가됨 (Added in version 2.0).

컨텍스트 매니저나 Dag() 생성자로 단일 Dag를 선언하는 더 전통적인 방법 외에도, 함수를 @dag로 데코레이션해 Dag 생성기 함수로 바꿀 수 있어요:

airflow/example_dags/example_dag_decorator.py [source]

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()

이 데코레이터는 Dags를 깔끔하게 만드는 새 방법일 뿐 아니라, 함수에 있는 어떤 파라미터든 Dag 파라미터로 설정해서 Dag를 트리거할 때 그 파라미터를 설정할 수 있게 해 줘요. 그런 다음 파라미터를 Python 코드에서, 또는 Jinja template 안의 {{ context.params }}로 접근할 수 있어요.

참고 (Note)

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

제어 흐름 (Control Flow)

기본적으로 Dag는 의존하는 모든 Tasks가 성공할 때만 Task를 실행해요. 하지만 이를 수정하는 여러 방법이 있어요:

  • Branching — 조건에 따라 어느 Task로 이동할지 선택
  • Trigger Rules — Dag가 태스크를 실행할 조건 설정
  • Setup and Teardown — setup/teardown 관계 정의
  • Latest Only — 현재에 대해 실행되는 Dags에서만 실행되는 특수한 형태의 branching
  • Depends On Past — 태스크가 이전 실행의 자신에 의존할 수 있음

Branching

branching을 사용해 Dag가 의존하는 모든 태스크를 실행하지 않고, 하나 이상의 경로를 골라가도록 지시할 수 있어요. 여기서 @task.branch 데코레이터가 등장해요.

@task.branch 데코레이터는 @task와 매우 유사하지만, 데코레이션된 함수가 태스크의 ID(또는 ID 목록)를 반환할 것으로 기대해요. 지정된 태스크는 수행되고, 다른 모든 경로는 skip돼요. 또한 None 을 반환해 모든 다운스트림 태스크를 skip할 수도 있어요.

Python 함수가 반환한 task_id는 @task.branch 데코레이션된 태스크의 직접 다운스트림에 있는 태스크를 참조해야 해요.

참고 (Note)

태스크가 branching operator 선택된 태스크 중 하나 이상의 다운스트림에 있으면 skip되지 않아요: branching 태스크의 경로는 branch_a, join, branch_b예요. joinbranch_a의 다운스트림 태스크이므로, branch 결정의 일부로 반환되지 않았어도 여전히 실행돼요.

@task.branch는 XCom과 함께 사용할 수도 있으며, 업스트림 태스크에 기반해 branching context가 동적으로 어떤 branch를 따를지 결정하게 할 수 있어요. 예를 들어:

@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]

branching 기능으로 자신만의 operator를 구현하고 싶다면 BaseBranchOperator에서 상속받을 수 있어요. 이는 @task.branch 데코레이터와 유사하게 동작하지만 choose_branch 메서드의 구현을 제공해야 해요.

참고 (Note)

@task.branch 데코레이터는 BranchPythonOperator의 Taskflow 동등물이에요. 후자는 일반적으로 커스텀 operator를 구현하기 위해서만 서브클래싱해야 해요.

@task.branch의 callable과 마찬가지로, 이 메서드는 실행될 다운스트림 태스크의 ID나 태스크 ID 목록을 반환할 수 있고, 나머지는 모두 skip돼요. 모든 다운스트림 태스크를 skip하려면 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이라는 branch 데코레이터도 있어요.

Latest Only

Airflow의 Dag Runs는 종종 현재 날짜와 같지 않은 날짜에 실행돼요 — 예를 들어 지난 달의 매일마다 Dag 사본 하나를 실행해 데이터를 backfill하는 경우예요.

하지만 이전 날짜의 Dag 실행의 일부(또는 전부)를 실행하지 않게 하고 싶은 상황이 있어요. 이 경우 LatestOnlyOperator를 사용할 수 있어요.

이 특별한 Operator는 현재 "최신" Dag run에 있지 않으면(지금의 벽시계 시간이 그것의 execution_time과 다음 스케줄 execution_time 사이이고, 외부 트리거된 run이 아닌 경우) 다운스트림의 모든 태스크를 skip해요.

예시는 다음과 같아요:

airflow/example_dags/example_latest_only_with_trigger.py [source]

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에 대해 skip돼요.
  • task2latest_only와 완전히 독립적이며 모든 스케줄 기간에 실행돼요.
  • task3task1task2의 다운스트림이고, 기본 trigger rule이 all_success이므로 task1로부터 연쇄 skip을 받아요.
  • task4task1task2의 다운스트림이지만 trigger_ruleall_done으로 설정되어 있으므로 skip되지 않아요.

Depends On Past

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

Dag를 수명의 맨 처음, 즉 첫 자동 실행에서 실행한다면 태스크는 여전히 실행된다는 점을 참고하세요. 의존할 이전 실행이 없기 때문이에요.

Trigger Rules

기본적으로 Airflow는 태스크의 모든 업스트림(직접 부모) 태스크가 성공할 때까지 기다렸다가 태스크를 실행해요.

하지만 이것은 기본 동작일 뿐이며, Task의 trigger_rule 인자로 제어할 수 있어요. trigger_rule의 옵션은 다음과 같아요:

  • all_success (기본값): 모든 업스트림 태스크가 성공함
  • all_failed: 모든 업스트림 태스크가 failed 또는 upstream_failed 상태
  • all_done: 모든 업스트림 태스크가 실행을 마침
  • all_done_setup_success: all_done과 같지만, 태스크에 업스트림 setup 태스크가 있으면 그중 적어도 하나가 성공해야 함. teardown 태스크의 기본 trigger rule.
  • all_done_min_one_success: 모든 비-skip 업스트림 태스크가 실행을 마치고 적어도 하나의 업스트림 태스크가 성공함
  • all_skipped: 모든 업스트림 태스크가 skipped 상태
  • one_failed: 적어도 하나의 업스트림 태스크가 실패함(모든 업스트림 태스크가 끝나기를 기다리지 않음)
  • one_success: 적어도 하나의 업스트림 태스크가 성공함(모든 업스트림 태스크가 끝나기를 기다리지 않음)
  • one_done: 적어도 하나의 업스트림 태스크가 성공하거나 실패함
  • none_failed: 모든 업스트림 태스크가 failedupstream_failed가 아님 — 즉 모든 업스트림 태스크가 성공하거나 skip되었음
  • none_failed_min_one_success: 모든 업스트림 태스크가 failedupstream_failed가 아니고, 적어도 하나의 업스트림 태스크가 성공함
  • none_skipped: 어떤 업스트림 태스크도 skipped 상태가 아님 — 즉 모든 업스트림 태스크가 success, failed, upstream_failed, 또는 removed 상태
  • always: 의존성이 전혀 없음, 어떤 때든 이 태스크를 실행

참고 (Note)

removed Task Instance 상태는 run이 시작된 이후 태스크가 Dag에서 사라졌다는 뜻이에요. trigger rule의 경우 removed는 터미널 상태로, all_done, all_done_setup_success, all_done_min_one_success 같은 규칙의 "done"에 포함되지만, success, failed, upstream_failed, 또는 skipped로는 계산되지 않아요. 동적으로 매핑된 태스크의 경우 removed 업스트림 map index도 all_success, all_failed, none_failed, none_failed_min_one_success, all_done_min_one_success의 실패 수에서 뺍니다.

원한다면 이것을 Depends On Past 기능과 결합할 수도 있어요.

참고 (Note)

trigger rule과 skip된 태스크 사이의 상호작용, 특히 branching 작업의 일부로 skip된 태스크를 인지하는 것이 중요해요. branching 작업의 다운스트림에서는 all_success나 all_failed를 거의 사용하고 싶지 않을 거예요.

Skip된 태스크는 trigger rule all_successall_failed를 통해 연쇄되며 그것들도 skip되게 해요. 다음 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로 설정되어 있고, branching 작업으로 인한 skip이 all_success로 표시된 태스크를 skip하도록 연쇄되기 때문에 skip된 것으로 표시돼요.

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

Setup and teardown

데이터 워크플로우에서는 리소스(예: 컴퓨트 리소스)를 만들고, 사용하고, 정리하는 것이 흔해요. Airflow는 이 필요를 지원하기 위해 setup/teardown 태스크를 제공해요.

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

동적 Dags (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 태스크의 토폴로지(배치)를 비교적 안정적으로 유지하는 것을 권장해요. 동적 Dags는 보통 구성 옵션을 동적으로 로드하거나 operator 옵션을 바꾸는 데 더 잘 사용돼요.

Dag 시각화 (Dag Visualization)

Dag의 시각적 표현을 보고 싶다면 두 가지 옵션이 있어요:

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

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

물론 Dags를 개발하면서 점점 복잡해지므로, Dag 뷰를 수정해 이해하기 쉽게 하는 몇 가지 방법을 제공해요.

TaskGroups

TaskGroup은 Graph 뷰에서 태스크를 계층적 그룹으로 구성하는 데 사용할 수 있어요. 반복 패턴을 만들고 시각적 혼란을 줄이는 데 유용해요.

TaskGroup의 태스크들은 같은 원래 Dag에 살며 모든 Dag 설정과 pool 구성을 존중해요.

더 보기 (See also)

TaskGrouptask_group의 API 참조

의존성 관계는 >><< 연산자로 TaskGroup의 모든 태스크에 적용할 수 있어요. 예를 들어 다음 코드는 task1task2를 TaskGroup group1에 넣고 두 태스크를 모두 task3의 업스트림에 둬요:

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로 프리픽스됩니다. 이는 Dag 전체에서 group_id와 task_id의 고유성을 보장하는 데 도움이 돼요.

프리픽스싱을 비활성화하려면 TaskGroup을 만들 때 prefix_group_id=False를 전달하세요. 다만 이제는 모든 개별 태스크와 그룹이 자체적으로 고유한 ID를 갖도록 보장할 책임이 있다는 점을 참고하세요.

참고 (Note)

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

Edge Labels

태스크를 그룹으로 묶는 것 외에도 Graph 뷰에서 서로 다른 태스크 사이의 의존성 엣지 에 라벨을 붙일 수 있어요. 이는 특히 Dag의 branching 영역에서 유용해서, 특정 branch가 실행될 수 있는 조건을 라벨로 붙일 수 있어요.

라벨을 추가하려면 >><< 연산자와 함께 인라인으로 직접 사용할 수 있어요:

from airflow.sdk import Label

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

또는 set_upstream/set_downstream에 Label 객체를 전달할 수 있어요:

from airflow.sdk import Label

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

다른 branch에 라벨을 붙이는 것을 보여주는 예시 Dag는 다음과 같아요:

airflow/example_dags/example_branch_labels.py [source]

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)

Dags, TaskGroups, task 객체에 웹 인터페이스에서 볼 수 있는 문서나 메모를 추가할 수 있어요.

정의되면 풍부한 콘텐츠로 렌더링되는 특수 task 속성 세트가 있어요:

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

Dags와 TaskGroups의 경우 doc_md만 해석되는 속성이라는 점을 참고하세요. 그것은 문자열이나 markdown 파일의 참조를 포함할 수 있어요. Markdown 파일은 .md로 끝나는 문자열로 인식돼요. 상대 경로가 제공되면 Airflow Scheduler나 Dag parser가 시작된 경로 기준으로 로드돼요. 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")

Dags 패키징 (Packaging Dags)

단순한 Dags는 보통 단일 Python 파일에 있지만, 더 복잡한 Dags가 여러 파일에 걸쳐 있고 함께 배송되어야 하는 의존성("vendored")을 가질 수 있는 것은 드문 일이 아니에요.

이 모든 것을 Dag bundle 안에서 표준 파일시스템 배치로 할 수도 있고, Dag와 모든 Python 파일을 단일 zip 파일로 패키징할 수도 있어요. 예를 들어 다음 내용의 zip 파일로 두 개의 Dags를 필요로 하는 의존성과 함께 배송할 수 있어요:

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

패키징된 Dags에는 몇 가지 함정이 있다는 점을 참고하세요:

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

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

.airflowignore

.airflowignore 파일은 Dag bundle이나 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])을 사용해 범위의 문자 중 하나와 일치할 수 있어요.
  • 패턴은 !를 프리픽스로 부정할 수 있어요. 패턴은 순서대로 평가되므로 부정이 같은 파일에서 이전에 정의된 패턴이나 부모 디렉토리에 정의된 패턴을 덮어쓸 수 있어요.
  • 이중 별표(**)를 사용해 디렉토리를 가로질러 일치할 수 있어요. 예를 들어 **/__pycache__/는 각 하위 디렉토리의 __pycache__ 디렉토리를 무한 깊이까지 무시해요.
  • 패턴의 시작이나 중간(또는 둘 다)에 /가 있으면 패턴은 특정 .airflowignore 파일 자체의 디렉토리 레벨에 상대적이에요. 그렇지 않으면 패턴은 .airflowignore 레벨 아래의 어떤 레벨에서도 일치할 수 있어요.

regexp 패턴 문법의 경우 .airflowignore의 각 줄은 정규식 패턴을 지정하고, 이름(즉 Dag id 아님)이 패턴과 일치하는 디렉토리나 파일은 무시돼요(내부적으로는 Pattern.search()가 패턴을 일치시키는 데 사용돼요). # 문자로 주석을 표시하세요. #로 시작하는 줄의 모든 문자는 무시돼요.

.airflowignore 파일은 Dag bundle에 두어야 해요. 예를 들어 glob 문법으로 .airflowignore 파일을 준비할 수 있어요:

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

그러면 Dag bundle의 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 bundle의 하위 폴더에 .airflowignore 파일을 준비할 수 있고 그것은 해당 하위 폴더에만 적용돼요.

Dag 의존성 (Dag Dependencies)

Airflow 2.1에서 추가됨

Dag 내 태스크 간의 의존성은 업스트림·다운스트림 관계로 명시적으로 정의되는 반면, Dag 간의 의존성은 조금 더 복잡해요. 일반적으로 한 Dag가 다른 Dag에 의존할 수 있는 두 가지 방법이 있어요:

추가적인 어려움은 한 Dag가 서로 다른 데이터 구간으로 다른 Dag의 여러 실행을 대기하거나 트리거할 수 있다는 것이에요. 이러한 의존성은 Dag 직렬화 중에 스케줄러가 계산해요.

의존성 감지기는 구성 가능하므로, DependencyDetector의 기본값과 다른 자신만의 로직을 구현할 수 있어요.

Dag pausing, 비활성화(deactivation)와 삭제 (Dag pausing, deactivation and deletion)

Dags에는 "실행되지 않는" 것과 관련된 몇 가지 상태가 있어요. Dags는 paused, deactivated될 수 있고, 마지막으로 Dag의 모든 메타데이터가 삭제될 수 있어요.

Dag가 DAGS_FOLDER에 있고 스케줄러가 데이터베이스에 저장했지만 사용자가 UI를 통해 비활성화하기로 선택했다면 UI를 통해 pause할 수 있어요. "pause"와 "unpause" 동작은 UI와 API를 통해 사용할 수 있어요. Paused Dags는 Scheduler가 스케줄하지 않지만, 수동 실행을 위해 UI로 트리거할 수 있어요. UI에서 paused Dags를 볼 수 있어요(Paused 탭). Un-pause된 Dags는 Active 탭에서 찾을 수 있어요. Dag가 pause되면 실행 중인 태스크는 완료가 허용되고 모든 다운스트림 태스크는 "Scheduled" 상태가 돼요. Dag가 unpause되면 "scheduled" 태스크는 Dag 로직에 따라 실행을 시작해요. "scheduled" 태스크가 없는 Dags는 스케줄에 따라 실행을 시작해요.

Dags는 DAGS_FOLDER에서 제거해 deactivate할 수 있어요(UI의 Active 탭과 혼동하지 마세요). 스케줄러가 DAGS_FOLDER를 파싱할 때 이전에 봤고 데이터베이스에 저장한 Dag를 찾지 못하면 그것을 deactivated로 설정해요. Deactivated Dags에 대해서는 Dag의 메타데이터와 이력이 보존되고, Dag가 DAGS_FOLDER에 다시 추가되면 다시 activated되고 이력이 보일 거예요. UI나 API로 Dag를 activate/deactivate할 수는 없고, DAGS_FOLDER에서 파일을 제거해서만 할 수 있어요. 다시 말하지만 스케줄러가 Dag를 deactivate할 때 Dag의 과거 실행에 대한 데이터는 손실되지 않아요. Airflow UI의 Active 탭은 Activated이면서 Not paused인 Dags를 가리킨다는 점을 참고하세요. 그래서 처음에는 조금 혼동될 수 있어요.

UI에서 deactivated Dags는 볼 수 없어요 — 가끔 과거 실행은 볼 수 있지만, 그 정보를 보려 하면 Dag가 누락되었다는 오류를 보게 될 거예요.

또한 UI나 API를 사용해 메타데이터 데이터베이스에서 Dag 메타데이터를 삭제할 수도 있지만, 그것이 항상 UI에서 Dag가 사라지는 결과를 주는 것은 아니에요 — 역시 처음에는 혼동될 수 있어요. 메타데이터를 삭제할 때 Dag가 여전히 DAGS_FOLDER에 있으면, Scheduler가 폴더를 파싱하며 Dag가 다시 나타날 거예요. Dag의 과거 실행 정보만 제거될 뿐이에요.

이 모든 것은 Dag와 그 모든 과거 메타데이터를 실제로 삭제하고 싶다면 세 단계로 해야 한다는 뜻이에요:

  • Dag를 pause한다
  • UI나 API로 데이터베이스에서 과거 메타데이터를 삭제한다
  • DAGS_FOLDER에서 Dag 파일을 삭제하고 그것이 비활성화될 때까지 기다린다

Dag 자동 pause (실험적) (Dag Auto-pausing (Experimental))

Dags는 자동으로 pause되도록 구성할 수도 있어요. 연속으로 N번 실패하면 Dag를 자동으로 비활성화하게 하는 Airflow 구성이 있어요.

또한 Dag 인자에서 이 구성을 제공하고 오버라이드할 수도 있어요:

Deadline Alerts

3.1 버전에 추가됨 (Added in version 3.1).

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

기존 email 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)

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

dag.test(use_executor=True)

pytest 테스트 스위트 안에서 호출을 사용하면 Airflow pytest 플러그인의 conf_vars fixture를 활용할 수 있는데, 이는 사전 정의된 구성 값의 쉬운 변경을 허용해요.

더 알아보기 (Learn more)