Airflow Listener Plugin
Airflow Listener Plugin
플러그인을 사용해 태스크 상태를 모니터링·추적하는 listener를 추가하는 Airflow 기능을 소개하는 문서예요. Task instance와 Dag run의 상태 변화를 감지하는 listener를 구현하는 방법을 예제 코드와 함께 살펴볼게요.
출처: 문서
본문
Airflow에는 플러그인을 사용해 listener를 추가해서 태스크 상태를 모니터링·추적하는 기능이 있어요.
이것은 Airflow의 간단한 예시 listener 플러그인으로, 태스크 상태를 추적하고 태스크·Dag run·dag에 대한 유용한 메타데이터를 수집하는 데 도움을 줘요.
이 플러그인은 SQLAlchemy의 이벤트 메커니즘을 이용해 동작해요. 테이블 레벨에서 태스크 인스턴스 상태 변경을 감시하고 이벤트를 트리거해요. 모든 Dag의 모든 태스크에 대해 알림이 전달돼요.
이 플러그인에서는 객체 참조가 기본 클래스 airflow.plugins_manager.AirflowPlugin에서 파생돼요.
Listener 플러그인은 내부적으로 pluggy 앱을 사용해요. Pluggy는 Pytest를 위해 만들어졌고 플러그인 관리와 hook 호출을 위한 앱이에요. Pluggy는 함수 hooking을 가능하게 해서, 그 hooking 위에 자신만의 커스터마이징으로 "pluggable" 시스템을 구축할 수 있어요.
이 플러그인을 사용하면 다음 이벤트를 들을 수 있어요:
- task instance가 running 상태일 때
- task instance가 success 상태일 때
- task instance가 failure 상태일 때
- Dag run이 running 상태일 때
- Dag run이 success 상태일 때
- Dag run이 failure 상태일 때
- Airflow job, scheduler 같은 이벤트 시작 전(on start before)
- Airflow job, scheduler 같은 이벤트 중지 전(before stop)
Listener 등록 (Listener Registration)
listener 객체에 대한 객체 참조가 있는 listener 플러그인이 Airflow 플러그인의 일부로 등록돼요. 다음은 새 listener를 구현하기 위한 골격(skeleton)이에요:
from airflow.plugins_manager import AirflowPlugin
# This is the listener file created where custom code to monitor is added over hookimpl
import listener
class MetadataCollectionPlugin(AirflowPlugin):
name = "MetadataCollectionPlugin"
listeners = [listener]
다음으로, listener에 추가된 코드를 확인하고 각 listener에 대한 구현 메서드를 볼 수 있어요. 구현 후 listener 부분은 모든 Dag의 모든 태스크 실행 중에 실행돼요.
참고로, 데이터베이스의 테이블 목록을 보여주는 listener.py 클래스 내부의 플러그인 코드는 다음과 같아요:
이 예시는 task instance가 running 상태일 때 이를 감지해요.
airflow/example_dags/plugins/event_listener.py [source]
@hookimpl
def on_task_instance_running(
previous_state: TaskInstanceState | None, task_instance: RuntimeTaskInstance | TaskInstance
):
"""
Called when task state changes to RUNNING.
previous_task_state and task_instance object can be used to retrieve more information about current
task_instance that is running, its dag_run, task and dag information.
"""
print("Task instance is in running state")
print(" Previous state of the Task instance:", previous_state)
name: str = task_instance.task_id
context = task_instance.get_template_context()
task = context["task"]
if TYPE_CHECKING:
assert task
dag = task.dag
dag_name = None
if dag:
dag_name = dag.dag_id
print(f"Current task name:{name}")
print(f"Dag name:{dag_name}")
마찬가지로, task_instance의 success·failure 이후를 감지하는 코드도 구현할 수 있어요.
이 예시는 Dag run이 failed 상태로 바뀔 때를 감지해요.
airflow/example_dags/plugins/event_listener.py [source]
@hookimpl
def on_dag_run_failed(dag_run: DagRun, msg: str):
"""
This method is called when dag run state changes to FAILED.
"""
print("Dag run in failure state")
dag_id = dag_run.dag_id
run_id = dag_run.run_id
run_type = dag_run.run_type
print(f"Dag information:{dag_id} Run id: {run_id} Run type: {run_type}")
print(f"Failed with message: {msg}")
마찬가지로, dag_run의 success 이후와 running 상태 동안을 감지하는 코드도 구현할 수 있어요.
listener 구현을 추가하는 데 필요한 listener 플러그인 파일은 Airflow 플러그인의 일부로 $AIRFLOW_HOME/plugins/ 폴더에 추가되고 Airflow 시작 시 로드돼요.