← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 33번째

Airflow 모듈 33/151 airflow-learn-33

Message Queues — 이벤트 기반 Dag 스케줄링

Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation core-concepts/message-queues.rst 전체

이 모듈을 다 읽으면

  • Airflow가 메시지 큐 이벤트로 Dag를 트리거하는 기본 메커니즘을 설명할 수 있다
  • Triggerer와 BaseTrigger가 이 기능에서 맡는 역할을 설명할 수 있다

Airflow는 시간·의존성 기반 스케줄링이 기본이지만, 메시지 큐 같은 외부 이벤트에 반응하는 네이티브 이벤트 기반 스케줄링도 지원한다. 이 모듈은 Triggerer가 폴링 기반으로 외부 메시지 큐를 감시해 Dag를 트리거하는 기본 메커니즘을 다룬다.

시간/의존성 기반 스케줄링을 넘어서

Apache Airflow는 원래 워크플로의 시간 기반·의존성 기반 스케줄링을 위해 설계되었다. 그러나 현대 데이터 아키텍처는 종종 준실시간(near real-time) 처리와, 메시지 큐 같은 다양한 소스에서 발생하는 이벤트에 반응하는 능력을 요구한다. Airflow는 네이티브 이벤트 기반 기능을 갖추고 있어, 외부 이벤트로 트리거되는 워크플로를 만들 수 있고 이를 통해 더 반응성 있는 데이터 파이프라인을 구축할 수 있다.

핵심 포인트

  • Airflow의 기본 스케줄링 모델은 시간/의존성 기반이지만, 이벤트 기반 스케줄링도 네이티브로 지원한다

폴링 기반 이벤트 스케줄링의 동작 방식

Airflow는 폴링 기반 이벤트 기반 스케줄링을 지원한다. Triggerer가 내장된 airflow.triggers.base.BaseTrigger 클래스를 이용해 외부 메시지 큐를 폴링할 수 있다. 이를 통해 사용자는 메시지 큐에 메시지가 도착하거나 데이터베이스의 변화 같은 외부 이벤트로 효율적으로 트리거되는 워크플로를 만들 수 있다.

Airflow는 외부 리소스의 상태를 지속적으로 모니터링하며, 그 리소스가 (도달한다면) 지정된 상태에 도달할 때마다 자산(asset)을 갱신한다. 이를 위해 Airflow Trigger를 활용하는데, Trigger는 외부 리소스 상태를 폴링하는 것이 임무인 작고 비동기적인 파이썬 코드 조각이다.

각 큐 provider는 provider별 키워드 인자를 받아 이를 내부의 Trigger에 그대로 전달한다. 지원되는 인자는 각 provider의 메시지 큐 문서를 참고해야 한다. Trigger가 Triggerer 안에서 실행되므로, 이런 인자로 참조되는 사용자 코드(예: Python dot-notation 문자열로 넘겨진 함수)는 반드시 Triggerer 프로세스에서 임포트 가능해야 한다. 지원되는 메시지 큐 목록은 provider 문서의 core-extensions/message-queues에서 관리된다.

핵심 포인트

  • Triggerer가 외부 메시지 큐를 폴링하며, 이 폴링 로직은 BaseTrigger를 상속한 작고 비동기적인 코드로 구현된다
  • 큐 provider별 kwargs로 전달되는 사용자 코드(문자열 dot-notation 함수 참조 등)는 Triggerer 프로세스에서 임포트 가능해야 한다
  • 지원되는 메시지 큐 목록은 provider별 문서에서 관리된다