콜백
콜백 (Callbacks)
이 페이지는 로깅·모니터링에서 중요한 요소인 Task 콜백을 다뤄요. 특정 Task가 실패했을 때 알림을 보내거나 DAG가 성공했을 때 콜백을 호출하는 식으로, Dag·Task의 상태 변화에 맞춰 동작을 실행할 수 있어요. 5가지 콜백 타입(on_success_callback, on_failure_callback, on_retry_callback, on_execute_callback, on_skipped_callback)과 예제를 알려드려요.
출처: 문서
본문
로깅·모니터링의 중요한 구성 요소는 Task 콜백을 사용해 주어진 DAG나 Task의 상태 변화에, 또는 주어진 DAG 안의 모든 Task에 걸쳐 조치를 취하는 것이에요. 예를 들어 특정 Task가 실패했을 때 알림을 보내거나, DAG가 성공했을 때 콜백을 호출하고 싶을 수 있어요.
콜백을 정의할 수 있는 세 곳이 있어요.
- DAG 정의에 설정된 콜백은 DAG 레벨에 적용돼요.
default_args를 사용하면 DAG 안의 각 Task에 콜백을 설정할 수 있어요.- Task 정의 자체 안에서 해당 콜백을 설정하면 개별 Task에 콜백을 설정할 수 있어요.
Note
콜백 함수는 worker가 실행해 DAG나 Task 상태가 변경될 때만 호출돼요. 따라서 명령 줄 인터페이스(CLI)나 사용자 인터페이스(UI)로 설정한 DAG·Task 변경은 콜백 함수를 실행하지 않아요.
Warning
콜백 함수는 Task가 완료된 뒤에 실행돼요. 콜백 함수에서의 에러는 Task 로그가 아니라 dag processor 로그에 나타나요. 기본적으로 dag processor 로그는 UI에 표시되지 않으며, 대신
$AIRFLOW_HOME/logs/dag_processor/latest/dags-folder/<the_path_for_your_dag>/DAG_FILE.py.log에서 찾을 수 있어요.
Note
Airflow 2.6.0부터 콜백은 콜백 함수의 리스트를 지원해요. 원하는 이벤트에 실행될 여러 함수를 지정할 수 있게 된 거죠. Dag/Task 콜백을 정의할 때 콜백 인자에 콜백 함수 리스트를 전달하면 돼요: 예를 들어
on_failure_callback=[callback_func_1, callback_func_2]
콜백 타입
특정 이벤트에서 트리거되는 다섯 가지 콜백 타입이 있어요:
| 이름 | 설명 | 사용 가능 위치 |
|---|---|---|
on_success_callback |
Dag가 성공했을 때 또는 Task가 성공했을 때 호출돼요. | Dag 또는 Task |
on_failure_callback |
Dag가 실패했을 때 또는 Task가 실패했을 때 호출돼요. | Dag 또는 Task |
on_retry_callback |
Task가 재시도 대상일 때 호출돼요. | Task |
on_execute_callback |
Task가 실행을 시작하기 직전에 호출돼요. | Task |
on_skipped_callback |
Task가 실행 중일 때 AirflowSkipException이 발생하면 호출돼요. 명시적으로 말하면, DAG 안의 선행 분기 결정이나 트리거 규칙 때문에 Task 실행이 스케줄링되지 않고 건너뛰어져 실행이 시작되지 않은 경우에는 호출되지 않아요. |
Task |
컨텍스트 매핑
Task 인스턴스의 런타임 정보를 담은 컨텍스트 매핑이 모든 콜백에 전달돼요. context에서 사용 가능한 변수의 전체 목록은 문서와 코드에 있어요.
Dag 콜백
컨텍스트 매핑은 Task 인스턴스의 실행을 설명하므로, Dag 콜백에 전달되는 컨텍스트에도 Task 인스턴스 변수가 포함되며, 선택되는 Task는 DAG의 상태에 따라 달라져요:
- 일반적인 실패 시에는, 가장 최근에 실패한 Task가 선택돼요.
- Dag run 타임아웃 시에는, 시작되었지만 끝나지 않은 가장 최근 Task가 전달돼요.
- Task가 교착 상태(deadlock)라면, 다음에 실행돼야 하지만 실행될 수 없었던 Task가 전달돼요.
- 성공 시에는, 가장 최근에 성공한 Task가 전달돼요.
Dag 콜백에서 Task 인스턴스 변수에 의존하는 것은 사람이 분석할 때를 제외하고는 권장하지 않아요. 그 변수들은 DAG 상태에 대한 일부 정보만 반영하니까요. 예를 들어 타임아웃은 여러 정체(stalling) Task에 의해 발생할 수 있지만, 컨텍스트에는 결국 하나만 선택돼 전달돼요.
Note
Airflow 3.2.0 이전에는 위 규칙이 적용되지 않았고, Dag 콜백에 전달되는 Task 인스턴스는 DAG 상태와 무관하게 DAG에서 사전식으로 가장 최근 Task가 선택됐어요.
예제
커스텀 콜백 메서드 사용하기
다음 예제에서 task1의 실패는 task_failure_alert 함수를 호출하고, DAG 레벨의 성공은 dag_success_alert 함수를 호출해요. 각 Task가 실행을 시작하기 전에는 task_execute_callback 함수가 호출돼요:
from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator
def task_execute_callback(context):
print(f"Task has begun execution, task_instance_key_str: {context['task_instance_key_str']}")
def task_failure_alert(context):
print(f"Task has failed, task_instance_key_str: {context['task_instance_key_str']}")
def dag_success_alert(context):
print(f"Dag has succeeded, run_id: {context['run_id']}")
with DAG(
dag_id="example_callback",
on_success_callback=dag_success_alert,
default_args={"on_execute_callback": task_execute_callback},
):
task1 = EmptyOperator(task_id="task1", on_failure_callback=[task_failure_alert])
task2 = EmptyOperator(task_id="task2")
task3 = EmptyOperator(task_id="task3")
task1 >> task2 >> task3
Notifier 사용하기
Dag 정의에서 on_*_callbacks에 인자로 전달해 Notifier를 사용할 수 있어요. 예를 들어 on_success_callback이나 on_failure_callback과 함께 사용해 Task나 Dag run의 상태에 따라 알림을 보낼 수 있어요.
커스텀 notifier 사용 예제는 다음과 같아요:
from airflow.sdk import DAG
from airflow.providers.standard.operators.bash import BashOperator
from myprovider.notifier import MyNotifier
with DAG(
dag_id="example_notifier",
on_success_callback=MyNotifier(message="Success!"),
on_failure_callback=MyNotifier(message="Failure!"),
):
task = BashOperator(
task_id="example_task",
bash_command="exit 1",
on_success_callback=MyNotifier(message="Task Succeeded!"),
)
커뮤니티가 관리하는 Notifier 목록은 Notifications를 참고해요. 커스텀 Notifier 작성에 대한 자세한 내용은 Notifiers how-to 페이지를 참고해요.
데드라인 알림 콜백
위의 Dag/Task 라이프사이클 콜백 외에도 Airflow는 데드라인 알림(Deadline Alert) 콜백을 지원해요. 이 콜백은 Dag run이 설정된 시간 임계값을 초과할 때 트리거돼요. 데드라인 알림 콜백은 AsyncCallback (Triggerer에서 실행) 또는 SyncCallback (executor에서 실행)을 사용하며, DAG의 deadline 파라미터로 구성돼요.
자세한 내용은 Deadline Alerts를 참고해요.