self.defer() 호출 방식과 Task Start/End 직결 최적화
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/deferring.rst - Triggering Deferral, Deferring multiple times, Triggering Deferral from Task Start, Exiting deferred task from Triggers (약 185-470줄)
이 모듈을 다 읽으면
- self.defer()의 인자와 재개 시 method_name이 받는 context/event의 의미를 설명할 수 있다
- start_from_trigger/end_from_trigger가 워커를 거치지 않고 어떻게 리소스를 아끼는지 설명할 수 있다
- execution_timeout이 defer 전후에 어떻게 누적 적용되는지 설명할 수 있다
self.defer(trigger, method_name, kwargs, timeout)는 TaskDeferred 예외를 일으켜 오퍼레이터를 트리거러로 넘기고, 재개 시 지정된 method_name이 context와 event를 받아 이어서 실행되며, start_from_trigger/end_from_trigger를 쓰면 워커를 거치지 않고 트리거러에서 태스크를 시작하거나 끝낼 수 있어 자원을 더 아낄 수 있다.
self.defer() 호출과 재개
오퍼레이터 안 어디서든 `self.defer(trigger, method_name, kwargs, timeout)`을 호출해 defer를 발생시킬 수 있다. `trigger`는 defer할 대상 트리거 인스턴스로 DB에 직렬화되어 저장되고, `method_name`은 재개 시 Airflow가 호출할 오퍼레이터의 메서드 이름이며, `kwargs`는 그 메서드에 추가로 넘길 키워드 인자(기본값 `{}`), `timeout`은 이 defer가 초과하면 태스크를 실패시킬 타임델타(기본값 `None`, 즉 타임아웃 없음)다.
재개될 때 Airflow는 `context`와 `event`를 자동으로 kwargs에 추가해 `method_name`에 전달하므로, 재개 메서드는 반드시 이 둘을 키워드 인자로 받도록 선언해야 한다. `event`에는 트리거가 오퍼레이터를 재개시킬 때 함께 보낸 payload가 담기며, 상태 코드나 결과를 가져올 URL처럼 유용한 정보일 수도 있고 그냥 무시해도 되는 정보(예: datetime)일 수도 있다.
첫 `execute()`든 이후의 `method_name` 메서드든 오퍼레이터가 거기서 return하면 태스크는 완료된 것으로 간주된다. `self.defer` 호출은 내부적으로 `TaskDeferred` 예외를 발생시키는 방식으로 동작하므로, `execute()` 안 어디서든 (중첩 호출 안에서도) 호출할 수 있고, 필요하면 같은 인자로 `TaskDeferred`를 직접 raise할 수도 있다.
`execution_timeout`은 defer 전후의 개별 실행이 아니라 태스크의 총 실행 시간(total runtime)을 기준으로 판정된다 - 즉 `execution_timeout`이 설정되어 있으면, 재개 후 단 몇 초만 실행했더라도 총 누적 시간이 초과되면 태스크가 실패할 수 있다.
핵심 포인트
- `self.defer(trigger, method_name, kwargs, timeout)`는 TaskDeferred 예외를 발생시키는 방식으로 동작하므로 execute() 안 어디서든(중첩 호출 안에서도) 호출할 수 있다
- 재개 시 Airflow가 `context`와 `event`를 kwargs에 자동으로 추가해 method_name에 전달하므로, 재개 메서드는 반드시 이 둘을 키워드 인자로 받아야 한다
- `execution_timeout`은 defer 전후 개별 실행이 아니라 태스크의 총 실행 시간을 기준으로 판정되므로, 재개 후 몇 초 만에 실패할 수도 있다
한 오퍼레이터에서 여러 번 defer하기
가변 길이 리스트를 순회하며 각 항목마다 defer하고 싶은 경우가 있다 - 예를 들어 여러 개의 쿼리를 DB에 순차 제출하거나 여러 파일을 하나씩 처리하는 상황이다.
이때 `method_name`을 `execute` 자기 자신으로 지정해 하나의 진입점만 갖도록 만들 수 있는데, 다만 이 경우 `execute`는 선택적 키워드 인자로 `event`도 받을 수 있어야 한다. 진행 상태(예: 현재까지 처리한 항목 인덱스)는 `self`에 저장해봐야 재개 시 사라지므로, `kwargs`를 통해 다음 `execute` 호출로 명시적으로 전달해야 한다.
핵심 포인트
- method_name을 execute 자기 자신으로 지정하면 하나의 진입점에서 반복적으로 defer→재개를 거치며 리스트의 다음 아이템을 처리하는 패턴을 만들 수 있다
- 이때 진행 상태(예: current_item_index)는 kwargs를 통해 다음 execute 호출로 명시적으로 전달해야 한다 - self에 저장해도 재개 시 사라진다
워커를 거치지 않는 시작/종료 (start_from_trigger / end_from_trigger)
2.10.0부터, 클래스 레벨 속성 `start_from_trigger=True`와 `start_trigger_args`(`StartTriggerArgs` 객체: `trigger_cls`, `trigger_kwargs`, `next_method`, `next_kwargs`, `timeout`)를 지정하면 태스크가 워커로 가지 않고 곧바로 트리거러로 defer될 수 있다. 이 값들은 인스턴스 레벨에서도 오버라이드할 수 있다.
동적 태스크 매핑과 함께 이 기능을 쓰려면 `__init__`에서 `start_from_trigger`와 `trigger_kwargs`를 반드시 정확히 같은 파라미터 이름으로 선언해야 한다 (다른 이름을 쓰면 이 기능이 동작하지 않는다). `start_from_trigger`가 True인 태스크를 매핑할 때는 `__init__` 메서드 전체가 스킵되며, 스케줄러는 `partial`/`expand`로 제공된 값(없으면 클래스 속성값으로 폴백)을 이용해 태스크를 executor로 보낼지 트리거러로 보낼지 결정한다 - 이 단계에서는 아직 XCom 값이 해석(resolve)되지 않는다.
트리거 실행이 끝나면 태스크는 워커로 다시 보내져 `next_method`를 실행하거나, 태스크 인스턴스가 그대로 끝날 수도 있다. 후자를 위한 것이 `end_from_trigger=True`다 - 이를 설정하면 트리거가 `TaskSuccessEvent`나 `TaskFailureEvent`를 yield해 워커로 돌아가지 않고 트리거러에서 곧바로 태스크 인스턴스를 성공/실패로 종료할 수 있다 (이 경우 태스크 상태와 필요하다면 XCom도 함께 설정된다). 트리거가 인스턴스를 직접 종료하는 경우 `method_name`은 의미가 없어져 `None`으로 둘 수 있다.
다만 트리거에서 직접 종료하는 기능은 리스너(plugin listener)가 그 Deferrable Operator에 통합되어 있지 않을 때만 동작한다 - `end_from_trigger=True`인 오퍼레이터에 리스너가 통합되어 있으면 Dag 파싱 시점에 예외가 발생해 이 제약을 알려준다.
핵심 포인트
- start_from_trigger=True와 클래스 속성 start_trigger_args(StartTriggerArgs)를 쓰면 태스크가 워커를 거치지 않고 바로 트리거러에서 시작되어, 새 워커를 띄우는 비용을 아낀다
- 동적 태스크 매핑과 함께 쓰려면 __init__에서 start_from_trigger와 trigger_kwargs를 정확히 같은 파라미터 이름으로 받아야 하고, 매핑 시에는 __init__ 전체가 스킵되며 스케줄러가 partial/expand 값(없으면 클래스 속성)으로 트리거러 제출 여부를 판단한다 - 이 단계에서 XCom 값은 아직 해석되지 않는다
- end_from_trigger=True이면 트리거가 TaskSuccessEvent/TaskFailureEvent를 yield해 워커로 돌아가지 않고 트리거러에서 바로 태스크 인스턴스를 종료할 수 있다 (이때 method_name은 의미 없음)
- 리스너(plugin listener)가 붙은 Deferrable Operator에는 end_from_trigger=True를 쓸 수 없다 - Dag 파싱 시 예외가 발생한다