Listeners — Pluggy 기반 이벤트 알림과 버전 호환성
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation administration-and-deployment/listeners.rst (전체)
이 모듈을 다 읽으면
- 리스너가 감지할 수 있는 다섯 가지 이벤트 범주(Lifecycle, DagRun, TaskInstance, Asset, Dag Import Error)를 구분할 수 있다
- 리스너 API가 '모든 Dag·오퍼레이터에 걸쳐' 동작한다는 제약과, 특정 Dag만 다루고 싶을 때의 대안(on_success_callback 등)을 설명할 수 있다
- Airflow 3에서 API를 통한 상태 변경도 리스너에 알림이 가도록 바뀐 지점을 설명할 수 있다
- 리스너 인터페이스의 하위 호환성이 한 방향으로만 보장된다는 점을 설명할 수 있다
리스너는 Pluggy 기반으로 Airflow의 여러 라이프사이클·상태 변경 이벤트를 구독해 알림을 받을 수 있게 해주는 고급 기능이다. Airflow 컴포넌트와 격리되어 있지 않아 신중하게 작성해야 하며, 특정 Dag만 골라 들을 수는 없다는 제약이 있다. 인터페이스는 시간이 지나며 진화해왔고(2.10.0의 error 필드, 3.0.0의 세션 인자 제거, 3.2.0의 on_task_instance_skipped 추가), 오래된 버전을 대상으로 작성된 리스너는 최신 Airflow와도 대체로 호환되지만 그 반대는 보장되지 않는다.
리스너가 감지하는 이벤트 범주
리스너를 작성하면 이벤트가 발생했을 때 Airflow가 알려주도록 만들 수 있다. 이 기능은 Pluggy가 뒷받침한다. 다만 리스너는 고급 기능으로, 이들이 실행되는 Airflow 컴포넌트로부터 격리되어 있지 않아 경우에 따라 인스턴스를 느리게 만들거나 심지어 다운시킬 수도 있으므로 작성 시 각별한 주의가 필요하다.
Airflow가 지원하는 이벤트는 다섯 범주다. Lifecycle Events(``on_starting``, ``before_stopping``)는 ``SchedulerJob`` 같은 Airflow ``Job``의 시작·정지 이벤트에 반응할 수 있게 해준다. DagRun State Change Events(``on_dag_run_running``, ``on_dag_run_success``, ``on_dag_run_failed``)는 ``DagRun``이 상태를 바꿀 때 발생하며, Airflow 3부터는 API를 통해 상태 변경이 트리거될 때도(``on_dag_run_success``, ``on_dag_run_failed``에 한해) 알림이 간다 — 예를 들어 UI에서 DagRun을 success로 표시하는 경우다. TaskInstance State Change Events(``on_task_instance_running``, ``on_task_instance_success``, ``on_task_instance_failed``, ``on_task_instance_skipped``)는 ``RuntimeTaskInstance``가 상태를 바꿀 때 발생해 ``LocalTaskJob`` 상태 변경에 반응할 수 있게 해준다 — Airflow 3부터는 ``on_task_instance_success``/``on_task_instance_failed``도 API를 통한 상태 변경(예: UI에서 태스크 인스턴스를 success로 표시) 시 알림이 가며, 이 경우 리스너는 ``RuntimeTaskInstance`` 대신 ``TaskInstance`` 인스턴스를 받는다. Asset Events(``on_asset_created``, ``on_asset_alias_created``, ``on_asset_changed``)는 Asset 관리 작업이 실행될 때 발생한다. Dag Import Error Events(``on_new_dag_import_error``, ``on_existing_dag_import_error``)는 Dag processor가 Dag 코드에서 임포트 오류를 발견해 메타데이터 DB 테이블을 갱신할 때 발생한다.
핵심 포인트
- 리스너는 Pluggy 기반이며 실행되는 Airflow 컴포넌트와 격리되어 있지 않아 인스턴스 안정성에 영향을 줄 수 있는 고급 기능이다
- Airflow 3부터 DagRun/TaskInstance의 success·failed 이벤트는 API를 통한 상태 변경(예: UI 조작)에도 리스너 알림이 간다
- API 트리거로 태스크 인스턴스 상태가 바뀐 경우 리스너는 RuntimeTaskInstance가 아니라 TaskInstance 인스턴스를 받는다
리스너 작성법과 적용 범위의 한계
리스너를 만들려면 ``airflow.listeners.hookimpl``을 임포트하고, 알림을 생성하고 싶은 이벤트에 대한 hookimpl을 구현하면 된다. Airflow는 hookspec으로 이 명세를 정의하며, 구현체는 hookspec에 정의된 것과 동일한 이름의 파라미터를 받아야 한다 — 그렇지 않으면 플러그인을 사용하려 할 때 Pluggy가 에러를 던진다. 다만 모든 메서드를 구현할 필요는 없으며, 많은 리스너가 한두 개 메서드만 구현한다. 리스너를 Airflow 설치에 포함하려면 Airflow Plugin의 일부로 넣으면 된다.
리스너 API는 모든 Dag와 모든 오퍼레이터에 걸쳐 호출되도록 설계되어 있다 — 즉 특정 Dag가 발생시키는 이벤트만 골라 들을 수는 없다. 그런 동작이 필요하다면 ``on_success_callback``이나 ``pre_execute`` 같은 방법을 쓰는 것이 맞다 — 이들은 특정 Dag 작성자나 오퍼레이터 제작자를 위한 콜백을 제공한다. 로그와 ``print()`` 호출은 리스너의 일부로 처리된다.
핵심 포인트
- 리스너는 airflow.listeners.hookimpl로 구현하며, hookspec과 다른 파라미터 이름을 쓰면 Pluggy가 에러를 던진다
- 리스너 API는 모든 Dag·오퍼레이터에 전역으로 걸리므로 특정 Dag만 골라 들을 수 없고, 그런 용도에는 on_success_callback·pre_execute를 대신 써야 한다
호환성과 인터페이스 변경 이력
리스너 인터페이스는 시간이 지나며 바뀔 수 있다. Pluggy 명세를 쓰기 때문에, 오래된 버전의 인터페이스를 대상으로 작성된 리스너 구현체는 대체로 더 새로운 버전의 Airflow와도 앞으로 호환(forward-compatible)된다. 하지만 그 반대는 보장되지 않는다 — 더 새로운 인터페이스 버전을 대상으로 구현된 리스너는 더 오래된 버전의 Airflow에서 동작하지 않을 수 있다. 단일 Airflow 버전만 타깃한다면 그 버전에 맞춰 구현을 조정하면 되므로 문제되지 않지만, 여러 버전의 Airflow와 함께 쓰일 수 있는 플러그인·확장을 작성한다면 이 점이 중요해진다.
예를 들어 ``on_task_instance_failed``에 2.10.0에서 추가된 ``error`` 필드처럼 인터페이스에 새 필드가 추가된 경우, 그 필드가 없는 이벤트 객체를 다루지 못하는 리스너 구현은 Airflow 2.10.0 이상에서만 동작한다. 여러 버전을 지원하려면 런타임에 ``importlib.metadata.version("apache-airflow")``으로 Airflow 버전을 확인하고, 신버전에서는 새 필드를 활용하는 구현을, 구버전에서는 그 필드 없이 동작하는 구현을 조건 분기로 나눠 등록하면 된다.
2.8.0에 리스너가 도입된 이후 인터페이스 변경 이력은 다음과 같다: 2.10.0에서는 ``on_task_instance_failed``에 ``error`` 필드가 추가되었다. 3.0.0에서는 ``on_task_instance_running``에서 ``session`` 인자가 제거되고 ``task_instance``가 ``RuntimeTaskInstance`` 인스턴스가 되었으며, ``on_task_instance_failed``·``on_task_instance_success``에서도 마찬가지로 ``session`` 인자가 제거되고 ``task_instance``가 worker에서는 ``RuntimeTaskInstance``, API server에서는 ``TaskInstance``가 되었다. 3.2.0에서는 ``on_task_instance_skipped``라는 새 리스너 메서드가 인터페이스에 추가되었다.
핵심 포인트
- 오래된 인터페이스 대상 리스너는 새 Airflow와 호환되기 쉽지만, 새 인터페이스 대상 리스너가 오래된 Airflow에서 동작한다는 보장은 없다 — 호환성은 한 방향으로만 보장된다
- 여러 Airflow 버전을 지원하려면 런타임에 버전을 확인해 신/구 인터페이스 구현을 조건 분기로 나눠 등록해야 한다
- 3.0.0에서 TaskInstance 계열 리스너의 session 인자가 제거되고 worker에서는 RuntimeTaskInstance, API server에서는 TaskInstance를 받도록 바뀌었으며, 3.2.0에서 on_task_instance_skipped가 새로 추가되었다