지연 가능 오퍼레이터

Airflow로 데이터 파이프라인 구축하기

Volker Janz

Senior Developer Advocate at Astronomer

센서 문제

워커 슬롯을 점유하는 센서

 

  • Poke 모드: 센서가 워커 슬롯을 계속 점유
  • 장시간 실행되는 센서가 다른 태스크를 굶김
Airflow로 데이터 파이프라인 구축하기

Reschedule 모드

Poke 모드 (기본값)

from airflow.providers.standard.sensors.filesystem \
    import FileSensor

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • 슬롯이 전체 시간 동안 점유됨

Reschedule 모드

from airflow.providers.standard.sensors.filesystem \
    import FileSensor

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • 확인 사이에 슬롯이 해제됨
Airflow로 데이터 파이프라인 구축하기

Deferrable 모드

Triggerer가 처리하는 지연 센서

  • 태스크가 Triggerer 프로세스에 제어를 넘김
  • Triggerer는 비동기 I/O로 수백 개의 대기를 동시 처리
  • 대기 중 워커 슬롯을 전혀 사용하지 않음
Airflow로 데이터 파이프라인 구축하기

Deferrable 활성화

from airflow.providers.standard.sensors.filesystem import FileSensor

wait_for_data = FileSensor(
    task_id="wait_for_data",
    filepath="/data/incoming/report.csv",

deferrable=True,
poke_interval=30, timeout=3600, )

 

  • poke_intervaltimeout은 계속 적용됨
Airflow로 데이터 파이프라인 구축하기

각 모드의 사용 시점

모드 적합한 경우 워커 슬롯
poke 짧은 대기 (1분 미만) 전체 시간 점유
reschedule 중간 대기, Triggerer 없음 확인 사이에 해제
deferrable 긴 대기, 다수 센서 사용 안 함

 

  • Triggerer가 실행 중이라면 deferrable이 거의 항상 최선의 기본값
Airflow로 데이터 파이프라인 구축하기

전역 deferrable 설정

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • 호환되는 모든 오퍼레이터가 자동으로 deferrable 모드를 사용
  • DAG 코드 변경 불필요
  • 필요 시 오퍼레이터별로 deferrable=False로 재정의 가능
Airflow로 데이터 파이프라인 구축하기

연습해 봅시다!

Airflow로 데이터 파이프라인 구축하기

Preparing Video For Download...