Dag 제어 흐름 — 분기, Trigger Rule, Depends On Past
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation core-concepts/dags.rst - Control Flow (Branching, Latest Only, Depends On Past, Trigger Rules, Setup and teardown)
이 모듈을 다 읽으면
- @task.branch가 다운스트림 태스크를 스킵시키는 방식과, 분기 이후 join 태스크가 예상과 다르게 스킵되는 이유를 설명할 수 있다
- trigger_rule의 각 옵션(all_success, all_done, none_failed_min_one_success 등)이 언제 태스크를 실행시키는지 설명할 수 있다
- removed 상태가 trigger_rule 판정에서 다른 터미널 상태와 어떻게 다르게 취급되는지 설명할 수 있다
Dag는 기본적으로 모든 업스트림 태스크가 성공해야 다음 태스크를 실행하지만, Branching·LatestOnly·Depends On Past·Trigger Rule 같은 메커니즘으로 이 기본 동작을 다양하게 바꿀 수 있다. 이 모듈은 그중에서도 분기와 트리거 규칙, 그리고 분기 이후 스킵이 예기치 않게 전파되는 함정을 다룬다.
Branching: @task.branch로 실행 경로 선택하기
브랜칭은 Dag가 의존하는 모든 태스크를 실행하는 대신, 하나 이상의 경로를 선택해서 진행하도록 지시하는 기능이다. ``@task.branch`` 데코레이터는 ``@task``와 거의 같지만, 데코레이트된 함수가 태스크 ID(또는 ID 목록)를 반환하도록 기대한다는 점이 다르다. 지정된 태스크는 계속 진행되고, 나머지 경로는 모두 스킵된다. ``None``을 반환하면 모든 다운스트림 태스크가 스킵된다. 함수가 반환하는 task_id는 반드시 ``@task.branch`` 태스크의 직접 다운스트림 태스크를 가리켜야 한다.
한 가지 주의할 점은, 어떤 태스크가 분기 오퍼레이터의 다운스트림이면서 동시에 선택된 경로들 중 하나의 다운스트림이기도 하다면, 그 태스크는 스킵되지 않는다는 것이다. 예를 들어 branch_a와 join과 branch_b라는 경로가 있고 join이 branch_a의 다운스트림이라면, join은 분기 결정에서 반환되지 않았더라도 여전히 실행된다.
XCom을 활용하면 업스트림 태스크의 결과에 따라 동적으로 분기를 결정할 수 있다. 직접 오퍼레이터를 구현하려면 ``BaseBranchOperator``를 상속하고 ``choose_branch`` 메서드를 구현하면 되며, 이는 ``@task.branch``와 유사하게 동작하지만 하나 이상의 다운스트림 task_id 목록 또는 None을 반환한다. ``@task.branch``는 전통적인 ``BranchPythonOperator``의 TaskFlow 대응 버전이며, 후자는 이제 일반적으로 커스텀 오퍼레이터를 만들 때의 서브클래싱 용도로만 쓰인다. 가상환경이나 외부 Python에서 분기 함수를 실행하는 ``@task.branch_virtualenv``, ``@task.branch_external_python`` 데코레이터도 있다.
핵심 포인트
- @task.branch는 다운스트림 태스크 ID(들)를 반환하는 함수로, 반환된 경로만 실행되고 나머지는 스킵되며 None을 반환하면 전체가 스킵된다
- 분기 오퍼레이터와 선택된 경로 양쪽 모두의 다운스트림에 있는 태스크(join 지점)는 분기 결정에 포함되지 않았더라도 스킵되지 않는다
- BaseBranchOperator를 상속해 choose_branch를 구현하면 직접 분기 오퍼레이터를 만들 수 있다
Trigger Rule: 업스트림 상태 조건을 세밀하게 제어하기
기본적으로 Airflow는 태스크를 실행하기 전에 모든 업스트림(직접 부모) 태스크가 성공(success) 상태가 될 때까지 기다린다. 이 동작은 태스크의 ``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``(스킵되지 않은 모든 업스트림이 완료되고 그중 최소 하나는 성공), ``all_skipped``(모든 업스트림이 skipped), ``one_failed``(업스트림 중 하나라도 실패하면, 나머지를 기다리지 않고 즉시), ``one_success``(업스트림 중 하나라도 성공하면 즉시), ``one_done``(업스트림 중 하나가 성공 또는 실패로 끝나면), ``none_failed``(모든 업스트림이 failed나 upstream_failed가 아님, 즉 성공했거나 스킵됨), ``none_failed_min_one_success``(none_failed 조건에 더해 최소 하나는 성공), ``none_skipped``(스킵된 업스트림이 하나도 없음, 즉 success/failed/upstream_failed/removed 중 하나), ``always``(의존성 없이 언제나 실행).
removed 상태(실행 도중 Dag에서 사라진 태스크)는 trigger rule 판정에서 특별하게 취급된다. all_done, all_done_setup_success, all_done_min_one_success 같은 규칙에서는 '완료(done)'로 집계되지만, success·failed·upstream_failed·skipped 어느 쪽으로도 집계되지 않는다. 동적으로 매핑된 태스크의 경우 removed된 업스트림 맵 인덱스는 all_success, all_failed, none_failed, none_failed_min_one_success, all_done_min_one_success의 실패 카운트에서도 제외된다.
핵심 포인트
- 기본 trigger_rule은 all_success이며, all_failed/all_done/one_failed/one_success/none_failed 등으로 업스트림 조건을 세밀하게 바꿀 수 있다
- all_done_setup_success는 teardown 태스크의 기본 trigger rule로, all_done 조건에 setup 태스크 중 최소 하나 성공 조건이 추가된 규칙이다
- removed 상태는 all_done류 규칙에서는 '완료'로 집계되지만 success/failed/upstream_failed/skipped 어느 것으로도 집계되지 않는 특수한 터미널 상태다
분기와 스킵 전파의 함정
trigger rule과 스킵된 태스크 사이의 상호작용에는 주의가 필요하다. 특히 분기 연산 다음에 오는 태스크에는 거의 항상 ``all_success``나 ``all_failed``를 쓰지 말아야 한다. 스킵된 태스크는 ``all_success``와 ``all_failed`` trigger rule을 타고 전파되어, 그 규칙을 가진 다운스트림 태스크도 함께 스킵시킨다.
예를 들어 branching → branch_a → follow_branch_a → join, 그리고 branching → branch_false → join으로 이어지는 Dag에서 branching이 branch_a를 선택했다면 branch_false는 스킵된다. join은 follow_branch_a와 branch_false 둘 다의 다운스트림인데, 기본 trigger rule인 all_success 때문에 branch_false의 스킵이 join까지 전파되어 join도 스킵된 것으로 표시된다. join의 trigger_rule을 ``none_failed_min_one_success``로 바꾸면, '실패한 업스트림이 없고 최소 하나는 성공'이라는 조건만 확인하므로 스킵된 branch_false가 있어도 join이 의도대로 실행된다.
또한 LatestOnly 연산(``LatestOnlyOperator``)은 현재 시각이 실행 시각과 다음 예약 실행 시각 사이가 아니거나 외부 트리거로 실행된 경우가 아니라면, 자신의 다운스트림에 있는 모든 태스크를 스킵시킨다는 점에서 분기의 특수한 형태로 볼 수 있다. Depends On Past는 ``depends_on_past=True``로 설정하면 태스크가 이전 Dag Run에서의 자기 자신의 이전 실행이 성공한 경우에만 실행되도록 하는데, Dag의 생애 첫 자동 실행에서는 참조할 이전 실행이 없으므로 그대로 실행된다.
핵심 포인트
- 스킵은 all_success/all_failed trigger rule을 타고 다운스트림으로 전파되므로, 분기 다음의 join 태스크에는 이 두 규칙을 쓰면 안 된다
- join 태스크의 trigger_rule을 none_failed_min_one_success로 바꾸면 분기로 인한 의도치 않은 스킵 전파를 막을 수 있다
- LatestOnlyOperator는 현재가 '최신' 실행이 아니면 다운스트림을 모두 스킵시키는 특수한 형태의 분기이며, depends_on_past=True인 태스크는 Dag의 첫 자동 실행에서는 이전 실행이 없으므로 그대로 실행된다