지연 가능 오퍼레이터와 트리거
지연 가능 오퍼레이터와 트리거 (Deferrable Operators & Triggers)
센서가 외부 상태를 기다리는 동안에도 worker 슬롯 하나를 통째로 점유하고 있으면, 단순한 대기 작업에 소중한 리소스가 낭비돼요. 센서의 reschedule 모드가 이 문제를 일부 해결하긴 하지만 고정 간격으로만 재실행할 수 있고, 시간이 아닌 다른 기준으로 재개하기는 어려워요. Airflow는 이런 경우를 위해 **지연 가능 오퍼레이터(deferrable operator)**와 **트리거(trigger)**를 제공해요.
본문
일반 오퍼레이터와 센서는 실행되는 동안, 심지어 대기 중에도 전체 worker 슬롯을 차지해요. 지연 가능 오퍼레이터는 대기가 필요한 지점에서 스스로 defer해서 worker 슬롯을 놓아주고, 그 대신 오퍼레이터가 필요로 하는 폴링·대기 작업을 트리거가 맡아요. 트리거가 폴링이나 대기를 끝내면 오퍼레이터를 재개하라는 신호를 보내요.
트리거는 단일 Python 프로세스에서 도는 작고 비동기적인 Python 코드 조각이에요. 비동기라서 여러 트리거가 Airflow의 triggerer 컴포넌트에서 효율적으로 공존할 수 있어요.
참고로 Airflow 3.2는 단일 worker 슬롯 안에서 동시 I/O를 수행하는 Python 네이티브 async 태스크도 지원해요. 지연 오퍼레이터는 외부 이벤트를 기다리는 동안 worker 슬롯을 놓아주는 반면, async 태스크는 태스크 프로세스를 유지한 채 공유 이벤트 루프로 작업을 다중화해요.
지연 메커니즘의 흐름은 이렇게 돼요.
- 실행 중인 오퍼레이터(태스크 인스턴스)가 다른 작업이나 조건을 기다려야 하는 지점에 도달하면, 재개해 줄 이벤트에 묶인 트리거와 함께 스스로 defer해서 worker를 다른 작업에 할당합니다.
- 새 트리거 인스턴스가 Airflow에 등록되고, triggerer 프로세스가 이를 집어 들습니다.
- 트리거가 발화될 때까지 실행되고, 발화되면 스케줄러가 원래 태스크를 재스케줄합니다.
- 스케줄러는 태스크를 worker 노드에서 재개하도록 큐에 넣습니다.
지연 가능 오퍼레이터 사용하기
Airflow가 미리 준비한 지연 가능 오퍼레이터(예: TimeSensor)를 쓰려면 두 가지만 하면 돼요.
- Airflow 설치에 일반
scheduler와 함께triggerer프로세스를 최소 1개 실행한다 - DAG에서 지연 가능 오퍼레이터/센서를 사용한다
기존 DAG을 지연 가능 오퍼레이터로 업그레이드할 때는 API 호환이 되는 센서 변형을 제공하니, DAG에 이 변형들을 넣으면 다른 변경 없이 지연 오퍼레이터를 쓸 수 있어요.
지연 가능 오퍼레이터 작성하기
지연 오퍼레이터를 직접 작성할 때 기억할 세 가지가 있어요.
- 오퍼레이터는 트리거와 함께 스스로 defer해야 해요. Airflow core에 포함된 트리거를 쓸 수도 있고 직접 작성할 수도 있어요.
- 지연되는 동안 오퍼레이터는 worker에서 중지·제거되고 상태가 자동으로 유지되지 않아요. 특정 메서드에서 재개하도록 지시하거나 특정 kwargs를 넘겨서 상태를 유지할 수 있어요.
- defer는 여러 번 할 수 있고, 오퍼레이터가 큰 작업을 하기 전이나 후에도, 그리고 특정 조건이 충족되면 할 수도 있어요.
가장 단순한 예시로 1시간을 기다렸다가 끝나는 센서를 만들어 볼게요. deferrable 옵션에 따라 지연하거나 그냥 잠자는 두 가지 경로를 제공해요.
import time
from datetime import timedelta
from typing import Any
from airflow.configuration import conf
from airflow.sdk import BaseSensorOperator, Context
from airflow.providers.standard.triggers.temporal import TimeDeltaTrigger
class WaitOneHourSensor(BaseSensorOperator):
def __init__(
self, deferrable: bool = conf.getboolean("operators", "default_deferrable", fallback=False), **kwargs
) -> None:
super().__init__(**kwargs)
self.deferrable = deferrable
def execute(self, context: Context) -> None:
if self.deferrable:
self.defer(
trigger=TimeDeltaTrigger(timedelta(hours=1)),
method_name="execute_complete",
)
else:
time.sleep(3600)
def execute_complete(
self,
context: Context,
event: dict[str, Any] | None = None,
) -> None:
# We have no more work to do here. Mark as complete.
return
트리거 작성하기
트리거는 BaseTrigger를 상속하는 클래스로 작성하고 세 메서드를 구현해요. __init__(트리거의 인자 저장), serialize(상태를 DB에 직렬화), run(async 코루틴 — 발화 시 TriggerEvent를 yield). 단순 예시를 볼게요.
import asyncio
from airflow.triggers.base import BaseTrigger, TriggerEvent
from airflow.sdk.timezone import utcnow
class DateTimeTrigger(BaseTrigger):
def __init__(self, moment):
super().__init__()
self.moment = moment
def serialize(self):
return ("airflow.providers.standard.triggers.temporal.DateTimeTrigger", {"moment": self.moment})
async def run(self):
while self.moment > utcnow():
await asyncio.sleep(1)
yield TriggerEvent(self.moment)
defer 호출로 지연 트리거하기
오퍼레이터 아무 데서나 self.defer(trigger, method_name, kwargs, timeout)를 호출하면 지연이 트리거돼요. 이 호출은 Airflow가 처리하는 특별한 예외를 발생시켜요. 인자는 다음과 같아요.
trigger: 지연 대상이 될 트리거 인스턴스. DB에 직렬화됩니다.method_name: 재개될 때 Airflow가 호출할 오퍼레이터의 메서드 이름.kwargs: (선택) 메서드 호출 시 넘길 추가 키워드 인자. 기본값은{}.timeout: (선택) 이 시간이 지나면 지연이 실패하고 태스크 인스턴스도 실패하는timedelta. 기본값None은 타임아웃 없음.
from __future__ import annotations
from datetime import timedelta
from typing import Any
from airflow.sdk import BaseSensorOperator, Context
from airflow.providers.standard.triggers.temporal import TimeDeltaTrigger
class WaitOneHourSensor(BaseSensorOperator):
def execute(self, context: Context) -> None:
self.defer(trigger=TimeDeltaTrigger(timedelta(hours=1)), method_name="execute_complete")
def execute_complete(self, context: Context, event: dict[str, Any] | None = None) -> None:
# We have no more work to do here. Mark as complete.
return
한 오퍼레이터가 항목 목록을 하나씩 처리하면서 각 항목 사이에 지연을 두는 패턴도 가능해요. 트리거의 run()이 TriggerEvent를 yield하면 오퍼레이터의 execute가 이벤트와 함께 다시 호출되고, 다음 항목을 처리하기 위해 다시 defer하는 식이에요.
import asyncio
from airflow.sdk import BaseOperator
from airflow.triggers.base import BaseTrigger, TriggerEvent
class MyItemTrigger(BaseTrigger):
def __init__(self, item):
super().__init__()
self.item = item
def serialize(self):
return (self.__class__.__module__ + "." + self.__class__.__name__, {"item": self.item})
async def run(self):
result = None
try:
# Somehow process the item to calculate the result
...
yield TriggerEvent({"result": result})
except Exception as e:
yield TriggerEvent({"error": str(e)})
class MyItemsOperator(BaseOperator):
def __init__(self, items, **kwargs):
super().__init__(**kwargs)
self.items = items
def execute(self, context, current_item_index=0, event=None):
last_result = None
if event is not None:
# execute method was deferred
if "error" in event:
raise Exception(event["error"])
last_result = event["result"]
current_item_index += 1
try:
current_item = self.items[current_item_index]
except IndexError:
return last_result
self.defer(
trigger=MyItemTrigger(item),
method_name="execute", # The trigger will call this same method again
kwargs={"current_item_index": current_item_index},
)
태스크 시작부터 지연하기 (StartTriggerArgs)
오퍼레이터의 __init__에서 start_trigger_args를 정의하고 start_from_trigger = True로 설정하면, 태스크 실행 시작 전에 바로 지연 상태로 들어가서 worker 슬롯을 전혀 쓰지 않고 트리거만 구동할 수 있어요. 인자는 다음과 같아요.
trigger_cls: 트리거 클래스의 import 가능한 경로.trigger_kwargs:trigger_cls초기화에 넘길 키워드 인자. 전부 Airflow가 직렬화할 수 있어야 한다는 게 이 기능의 주요 제약이에요.next_method: 재개될 때 호출할 오퍼레이터의 메서드 이름.next_kwargs:next_method호출 시 넘길 추가 키워드 인자.timeout: (선택) 지연이 실패하는 타임아웃timedelta. 기본값None은 타임아웃 없음.
from __future__ import annotations
from datetime import timedelta
from typing import Any
from airflow.sdk import BaseSensorOperator, Context, StartTriggerArgs
class WaitOneHourSensor(BaseSensorOperator):
start_trigger_args = StartTriggerArgs(
trigger_cls="airflow.providers.standard.triggers.temporal.TimeDeltaTrigger",
trigger_kwargs={"moment": timedelta(hours=1)},
next_method="execute_complete",
next_kwargs=None,
timeout=None,
)
start_from_trigger = True
def execute_complete(self, context: Context, event: dict[str, Any] | None = None) -> None:
# We have no more work to do here. Mark as complete.
return
동적 태스크 매핑 지원을 켜려면 __init__ 메서드에서 start_from_trigger와 trigger_kwargs를 정의해야 해요. 그렇지 않으면 매핑된 모든 인스턴스가 같은 지연 설정을 공유하게 되거든요. 매핑 인스턴스마다 지연 시간을 다르게 주고 싶다면 이렇게 __init__에서 받아서 설정해요.
from __future__ import annotations
from datetime import timedelta
from typing import Any
from airflow.sdk import BaseSensorOperator, Context, StartTriggerArgs
class WaitHoursSensor(BaseSensorOperator):
start_trigger_args = StartTriggerArgs(
trigger_cls="airflow.providers.standard.triggers.temporal.TimeDeltaTrigger",
trigger_kwargs={"moment": timedelta(hours=1)},
next_method="execute_complete",
next_kwargs=None,
timeout=None,
)
start_from_trigger = True
def __init__(
self,
*args: list[Any],
trigger_kwargs: dict[str, Any] | None,
start_from_trigger: bool,
**kwargs: dict[str, Any],
) -> None:
# This whole method will be skipped during dynamic task mapping.
super().__init__(*args, **kwargs)
self.start_trigger_args.trigger_kwargs = trigger_kwargs
self.start_from_trigger = start_from_trigger
def execute_complete(self, context: Context, event: dict[str, Any] | None = None) -> None:
# We have no more work to do here. Mark as complete.
return
이제 WaitHoursSensor를 매핑해서 인스턴스마다 다른 대기 시간을 줄 수 있어요.
WaitHoursSensor.partial(task_id="wait_for_n_hours", start_from_trigger=True).expand(
trigger_kwargs=[{"hours": 1}, {"hours": 2}]
)
트리거에서 지연 태스크 종료하기
기본적으로 트리거가 발화하면 오퍼레이터의 재개 메서드가 호출돼요. 그런데 어떤 경우(예: 5시간짜리 비동기 센서)는 트리거가 태스크를 직접 종료하는 게 더 자연스러울 수 있어요. 이때는 method_name=None으로 defer하고 트리거가 TaskSuccessEvent를 yield하게 하면 되죠.
class WaitFiveHourSensorAsync(BaseSensorOperator):
# this sensor always exits from trigger.
def __init__(self, **kwargs) -> None:
super().__init__(**kwargs)
self.end_from_trigger = True
def execute(self, context: Context) -> NoReturn:
self.defer(
method_name=None,
trigger=WaitFiveHourTrigger(duration=timedelta(hours=5), end_from_trigger=self.end_from_trigger),
)
class WaitFiveHourTrigger(BaseTrigger):
def __init__(self, duration: timedelta, *, end_from_trigger: bool = False):
super().__init__()
self.duration = duration
self.end_from_trigger = end_from_trigger
def serialize(self) -> tuple[str, dict[str, Any]]:
return (
"your_module.WaitFiveHourTrigger",
{"duration": self.duration, "end_from_trigger": self.end_from_trigger},
)
async def run(self) -> AsyncIterator[TriggerEvent]:
await asyncio.sleep(self.duration.total_seconds())
if self.end_from_trigger:
yield TaskSuccessEvent()
else:
yield TriggerEvent({"duration": self.duration})
위 예시에서 트리거는 end_from_trigger가 True면 TaskSuccessEvent를 yield해서 태스크 인스턴스를 바로 종료하고, 아니면 오퍼레이터에 지정된 메서드로 태스크를 재개해요.
고가용성 (High Availability)
Airflow는 트리거를 한 곳에서만 동시에 실행하려 하고, 실행 중인 모든 triggerer에 하트비트를 유지해요. triggerer가 죽거나 Airflow DB가 있는 네트워크에서 분리되면, 그 호스트에 있던 트리거를 자동으로 다른 곳에 재스케줄해요. 트리거를 재스케줄하기 전에 Airflow는 머신이 복귀할 때까지 2.1 * triggerer.job_heartbeat_sec 초를 기다려요.
HA triggerer의 워크로드를 분배하는 설정도 있어요. max_trigger_to_select_per_loop의 기본값은 50, capacity의 기본값은 1000이에요.
트리거별로 triggerer 호스트 배정을 제어할 수도 있어요. 멀티팀 모드를 쓰면 --team-name 옵션이 모든 트리거 유형(태스크 생성, 이벤트 기반, 콜백)에 네이티브 팀 범위 트리거 배정을 제공해요. 예를 들어 두 개의 triggerer 호스트를 큐로 나눠 띄울 수 있어요.
# triggerer "X" startup command
airflow triggerer --queues=alice,bob
# triggerer "Y" startup command
airflow triggerer --queues=test_q
센서의 mode='reschedule'과 deferrable=True의 차이
센서는 하위 태스크로 진행하기 전에 특정 조건이 충족되기를 기다려요. 대기 시간을 관리하는 옵션으로 mode='reschedule'과 deferrable=True 두 가지가 있어요. mode='reschedule'은 조건이 충족될 때까지 스스로 계속 재스케줄하는 방식이고, deferrable=True는 유휴 시 실행을 일시 정지했다가 조건이 변하면 재개하는 방식이에요. 참고로 deferrable=True는 일부 오퍼레이터가 "나중에 재시도(재연기) 가능"을 나타내는 규약일 뿐, Airflow의 내장 파라미터나 모드가 아니에요. 대기 유형(시간 기반 vs 외부 변경 기반)에 따라 선택이 달라져요.