콜백: Dag/Task 상태 변화에 대응하기
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation administration-and-deployment/logging-monitoring/callbacks.rst (전체)
이 모듈을 다 읽으면
- 5가지 콜백 타입과 각각이 Dag/Task 레벨 중 어디에 적용 가능한지 구분할 수 있다
- Dag 콜백에 전달되는 태스크 인스턴스 선택 규칙을 설명할 수 있다
- 커스텀 콜백 함수와 Notifier 방식의 차이, 그리고 Deadline Alert 콜백과의 차이를 설명할 수 있다
콜백은 Dag나 태스크의 상태 변화에 반응해 알림을 보내거나 후처리를 하는 로깅/모니터링의 핵심 구성요소다. 정의 위치(Dag/default_args/개별 태스크), 실행 시점과 조건, 그리고 Dag 콜백에 어떤 태스크 인스턴스가 컨텍스트로 전달되는지의 규칙을 정확히 아는 것이 중요하다.
콜백 정의 위치와 실행 조건
콜백은 세 곳에서 정의할 수 있다. Dag 정의에 설정한 콜백은 Dag 레벨에 적용되고, ``default_args``를 통해서는 Dag 안 각 태스크에 콜백을 설정할 수 있으며, 태스크 정의 자체에서 개별 콜백을 지정할 수도 있다.
콜백 함수는 워커에 의한 실행으로 Dag나 태스크 상태가 바뀔 때만 호출된다. 즉 CLI나 UI로 상태를 바꾼 경우에는 콜백이 실행되지 않는다. 또한 콜백 함수는 태스크가 완료된 이후에 실행되며, 콜백 함수 안에서 발생한 에러는 태스크 로그가 아니라 dag processor 로그에 나타난다 — 이 로그는 기본적으로 UI에 보이지 않고 ``$AIRFLOW_HOME/logs/dag_processor/latest/dags-folder/<dag 경로>/DAG_FILE.py.log``에서 찾을 수 있다.
2.6.0부터 콜백은 함수 리스트를 지원해, ``on_failure_callback=[callback_func_1, callback_func_2]``처럼 여러 함수를 한 이벤트에 등록할 수 있다.
핵심 포인트
- 콜백은 Dag 정의, default_args, 개별 태스크 정의 세 곳에서 설정할 수 있다
- 워커 실행으로 인한 상태 변화에만 콜백이 호출되며, CLI/UI로 바꾼 상태 변화는 콜백을 트리거하지 않는다
- 콜백 함수의 에러는 태스크 로그가 아니라 dag processor 로그에 남고 기본적으로 UI에 보이지 않는다
- 2.6.0부터 하나의 이벤트에 콜백 함수 리스트를 등록할 수 있다
5가지 콜백 타입
콜백은 다섯 가지 이벤트에 대해 트리거된다. ``on_success_callback``은 Dag 또는 태스크가 성공했을 때(Dag/Task 모두 가능), ``on_failure_callback``은 Dag 또는 태스크가 실패했을 때(Dag/Task 모두 가능) 호출된다. ``on_retry_callback``은 태스크가 재시도 대기 상태가 될 때(태스크 전용), ``on_execute_callback``은 태스크 실행이 시작되기 직전(태스크 전용)에 호출된다. ``on_skipped_callback``은 태스크가 실행 중 ``AirflowSkipException``을 발생시켰을 때(태스크 전용) 호출되며, 분기(branching) 결정이나 트리거 규칙 때문에 애초에 태스크 실행이 스케줄되지 않아 건너뛰어진 경우에는 명시적으로 호출되지 않는다.
핵심 포인트
- on_success_callback, on_failure_callback은 Dag/Task 모두에 설정 가능하다
- on_retry_callback, on_execute_callback, on_skipped_callback은 Task 전용이다
- on_skipped_callback은 AirflowSkipException으로 인한 스킵에만 호출되고, 분기/트리거 규칙으로 아예 스케줄되지 않은 스킵에는 호출되지 않는다
컨텍스트 매핑과 Dag 콜백의 태스크 선택 규칙
모든 콜백에는 태스크 인스턴스에 대한 런타임 정보를 담은 컨텍스트 매핑이 전달된다. Dag 콜백의 경우 컨텍스트는 Dag의 상태에 따라 선택되는 태스크 인스턴스 변수를 포함한다: 일반 실패 시에는 가장 최근에 실패한 태스크, Dag run 타임아웃 시에는 가장 최근에 시작했지만 끝나지 않은 태스크, 데드락 상태라면 다음에 실행되어야 했지만 못한 태스크, 성공 시에는 가장 최근에 성공한 태스크가 전달된다.
Dag 콜백에서 태스크 인스턴스 변수는 사람이 분석하는 용도 외에는 신뢰하지 않는 것이 권장된다 — 이는 Dag 상태에 대한 부분적인 정보만 반영하기 때문이다. 예를 들어 타임아웃은 여러 태스크가 지연되어 발생할 수 있지만 컨텍스트에는 그중 하나만 선택되어 담긴다.
Airflow 3.2.0 이전에는 이 규칙이 적용되지 않았고, Dag 콜백에 전달되는 태스크 인스턴스는 Dag 상태와 무관하게 사전순(lexicographically)으로 가장 마지막인 태스크가 선택되었다 — 이는 버전에 따라 동작이 달라지는 부분이므로 어느 버전을 기준으로 하는지 명시할 필요가 있다.
핵심 포인트
- Dag 콜백은 상태별로 다른 태스크가 선택된다: 실패→최근 실패 태스크, 타임아웃→최근 시작 미완료 태스크, 데드락→다음 실행 예정이었으나 못한 태스크, 성공→최근 성공 태스크
- Dag 콜백의 태스크 인스턴스 변수는 부분 정보이므로 사람이 보는 용도 외에는 의존하지 않는 게 좋다
- (버전 민감) 3.2.0 이전에는 Dag 상태와 무관하게 사전순 마지막 태스크가 전달되었다
커스텀 콜백 함수, Notifier, Deadline Alert 콜백
가장 단순한 방식은 파이썬 함수를 정의해 ``on_failure_callback``, ``on_success_callback`` 등에 넘기는 것이다. 이와 달리 Notifier는 재사용 가능한 객체로, ``on_success_callback=MyNotifier(message="Success!")``처럼 인스턴스를 직접 콜백 인자로 전달한다. 커뮤니티가 관리하는 Notifier 목록은 프로바이더 문서에서 확인할 수 있다.
Dag/태스크 생명주기 콜백과 별개로 Airflow는 **Deadline Alert** 콜백을 지원한다. 이는 Dag run이 설정된 시간 임계값을 초과했을 때 트리거되며, ``AsyncCallback``(트리거러에서 실행)이나 ``SyncCallback``(executor에서 실행)을 사용하고 Dag의 ``deadline`` 파라미터로 구성한다.
핵심 포인트
- Notifier는 함수 대신 재사용 가능한 객체를 on_*_callback에 직접 전달하는 방식이다
- Deadline Alert 콜백은 생명주기 콜백과 별개로, Dag run이 시간 임계값을 넘었을 때 트리거된다
- Deadline Alert는 AsyncCallback(트리거러) 또는 SyncCallback(executor)으로 실행되며 deadline 파라미터로 설정한다