on_failure_callback를 통한 Dag 수준 재시도
on_failure_callback를 통한 Dag 수준 재시도 (Dag-level Retry via on_failure_callback)
태스크가 재시도를 모두 소진한 뒤 실패했을 때, Dag 실행 전체를 다시 시도하는 패턴을 설명하는 문서예요. on_failure_callback과 공개 REST API의 Clear 엔드포인트를 조합해 실패한 Dag run을 지우고 스케줄러가 다시 실행하게 하는 레시피예요. 재시도 횟수 제한, 한계(caveats), 자동 복구 가드까지 함께 살펴볼게요.
출처: 문서
본문
Airflow의 내장 재시도 메커니즘은 task 수준에서 동작해요: 각 태스크는 자신만의 retries와 retry_delay를 갖고, 실패 시 그 단일 태스크만 다시 실행돼요.
태스크는 이상적으로 task 수준에서 멱등(idempotent) 하게 설계되어 내장 태스크별 재시도 메커니즘으로 충분해야 해요. 이것이 권장되는 출발점이고, 대부분의 워크플로우는 그것 이상이 필요하지 않아요.
이 how-to는 태스크 수준 멱등성으로도 더 큰 단위의 재시도 — 사실상 Dag 수준 재시도 — 를 원하는 경우를 위한 것이에요. 예를 들어 중간 상태를 부분적으로 정리하기 어려운 다단계 파이프라인은, 실패한 단일 태스크를 다시 실행하는 것보다 처음부터 다시 실행하는 게 더 안전할 수 있어요.
이 문서는 on_failure_callback과 공개 Airflow REST API를 결합해 Dag 수준 재시도에 근사하는 패턴을 보여줘요: 태스크가 실패하면 콜백이 Clear 엔드포인트를 호출해 실패한 Dag run을 지우고, 스케줄러가 다시 실행하게 해요.
이것은 기존 프리미티브 위에 구축된 레시피이지 새 기능이 아니에요. 의견이 반영되어 있고 트레이드오프가 있어요. 사용하기 전에 아래 '한계(Caveats)'를 주의 깊게 읽어 보세요.
이에 대한 더 일반적인 구조 — 때때로 transactional task group 이라고도 불려요 — 가 Airflow dev list에서 추가 기능으로 논의되고 있어요. 유스케이스에 도움이 된다면, 이 레시피를 장기적으로 의존하기보다 그 논의에 참여하세요.
이 패턴이 언제 적합할까요?
이런 경우에 사용하세요:
- 작업의 단위가 자연스럽게 단일 태스크가 아니라 Dag run일 때.
- 태스크 재실행이 멱등할 때 — 다시 실행해도 중복된 부작용, 이중 청구, 불일치한 외부 상태가 생기지 않을 때.
- 시간과 비용 측면에서 제한된 횟수의 전체 재실행이 허용 가능할 때.
이런 경우엔 피하세요:
- 태스크에 비멱등 부작용(이메일 보내기, 결제 청구, 중복 제거되지 않는 API에 게시)이 있고 멱등하게 만들지 않았을 때.
- 재시도가 다운스트림 Dag나 asset에 투명해야 할 때 — Dag run을 지우면 원래 실행의 재시도가 아니라 새로운 시도가 생겨요.
- 실제로 보이는 실패 모드를 단순한 태스크별
retries설정이 이미 커버할 때.
동작 방식 (How it works)
- 태스크가 태스크 수준 재시도를 모두 소진한 뒤 실패해요.
- 태스크의
on_failure_callback이 dag processor에서 실행돼요. - 콜백이 (자신이 유지하는 카운터를 사용해 — 아래 '재시도 횟수 제한' 참고) 또 다른 Dag 수준 시도가 허용되는지 판단해요.
- 허용된다면 콜백이 Airflow REST API를 호출해 실패한 Dag run을 지워요.
- 스케줄러가 지워진 task instance들을 집어 다시 실행해요.
사용되는 REST 엔드포인트는 POST /api/v2/dags/{dag_id}/dagRuns/{dag_run_id}/clear예요.
전제 조건 (Prerequisites)
- dag processor에서 Airflow REST API로의 인증된 접근, 대상 Dag의 Dag run을 지울 권한.
- dag processor에서 API 서버로의 네트워크 도달성.
참고 (Note)
on_failure_callback은 정상적인 태스크 실행을 통해 실패가 발생할 때만 실행돼요. 수동으로 수행한 상태 변경(UI, CLI)은 이 콜백을 호출하지 않아요.
기본 예시: 실패한 Dag run 지우기 (Basic example: clear the failed Dag run)
다음 예시는 태스크 수준 재시도를 소진한 뒤 Dag run 내의 어떤 태스크든 실패하면 전체 Dag run을 지워요. 콜백이 첫 실패에 실행되고 전체 Dag run이 재시도 책임을 넘겨받도록 태스크에 retries=0을 설정했어요.
이 최소 버전은 실패 시 항상 지워요. 프로덕션에서 그대로 실행하지 마세요 — 태스크가 계속 실패하면 무한 루프에 빠질 수 있어요. 배포 전에 시도 횟수 제한을 추가하세요. 아래 '재시도 횟수 제한'을 참고하세요.
from urllib.parse import quote
import requests
from airflow.sdk import DAG, task
AIRFLOW_API_BASE = "https://airflow.example.com/api/v2" # your deployment's API base
def clear_dag_run_on_failure(context):
dag_run = context["dag_run"]
dag_id_path = quote(dag_run.dag_id, safe="")
run_id_path = quote(dag_run.run_id, safe="")
response = requests.post(
f"{AIRFLOW_API_BASE}/dags/{dag_id_path}/dagRuns/{run_id_path}/clear",
headers={"Authorization": "Bearer <token>"},
json={"dry_run": False},
timeout=30,
)
response.raise_for_status()
with DAG(
dag_id="example_dag_level_retry",
default_args={
"retries": 0, # let the Dag-level retry take over
"on_failure_callback": clear_dag_run_on_failure,
},
):
@task
def step_one(): ...
@task
def step_two(value): ...
step_two(step_one())
재시도 횟수 제한하기 (Limiting the number of retries)
항상 run을 지우는 순진한 콜백은 태스크가 계속 실패하면 무한 루프를 돌 거예요. dag_run.conf는 run이 생성될 때 설정되고 clear에 의해 변경되지 않으므로 재시도 카운터로 사용할 수 없어요 — 시도 횟수를 직접 추적해야 해요. 간단한 방법은 run에 스코프된 Variable을 사용하는 거예요:
from airflow.sdk import Variable
def _attempts_key(dag_id: str, run_id: str) -> str:
return f"dag_run_attempts::{dag_id}::{run_id}"
def _read_attempts(dag_id: str, run_id: str) -> int:
return int(Variable.get(_attempts_key(dag_id, run_id), default=0))
def _bump_attempts(dag_id: str, run_id: str) -> int:
new_value = _read_attempts(dag_id, run_id) + 1
Variable.set(_attempts_key(dag_id, run_id), str(new_value))
return new_value
지우기로 결정하기 전에 _read_attempts를 사용하고, API 호출 직전에 _bump_attempts를 호출해요. Dag run이 끝날 때(성공 시와 포기 시 모두) 카운터를 리셋해서 Variables가 쌓이지 않게 해요:
def cleanup_attempts(context):
dag_run = context["dag_run"]
Variable.delete(_attempts_key(dag_run.dag_id, dag_run.run_id))
cleanup_attempts를 Dag 수준의 on_success_callback에 연결하고, on_failure_callback의 포기(give-up) 분기에서 반환 전에 호출해요.
참고 (Note)
이 Variable 기반 카운터는 동시 콜백 간에 원자적(atomic)이지 않아요. Dag에 동시에 실패할 수 있는 병렬 분기가 있다면, 두 콜백이 같은 값을 읽고 둘 다 지우기로 결정할 수 있어요. 실제로는 두 번째 clear가 no-op(first가 이미 run을 리셋)이지만 중복 API 호출이 보일 수 있어요. 더 엄격한 보증이 필요하다면 compare-and-swap 시맨틱이 있는 백엔드를 사용하세요.
자동 복구 Dags 가드하기 (Guarding automated remediation Dags)
Dag 수준 재시도 패턴은 실패한 작업을 지우거나, 메시지를 재생하거나, 외부 job을 재시작하는 등 외부 시스템에 영향을 주는 작업을 실행하는 자동 복구 Dags에서도 사용할 수 있어요. 이러한 작업은 지속적인 실패가 Dag 실행을 가로질러 같은 복구 작업을 반복 실행하게 하지 않도록 가드레일(guardrail)이 있어야 해요.
Pool은 동시에 실행되는 복구 태스크 수를 제한할 수 있지만, 같은 복구 작업이 시도되는 빈도는 제한하지 않아요. 자동 복구 워크플로우에서는 각 대상에 작은 cooldown 마커와 수동 검토 전 최대 자동 작업 횟수를 추가하는 것을 고려해 보세요.
cooldown 타임스탬프나 시도 카운터 같은 작은 가드 값에는 Airflow Variable로 충분할 수 있어요. Variables를 대용량 상태 저장소나 강한 일관성을 가진 잠금 메커니즘으로 사용하지 마세요. 워크플로우가 원자적 업데이트, 엄격한 동시성 제어, 또는 대상별 상태 레코드가 많아야 한다면 그 목적에 맞게 설계된 외부 저장소를 사용하세요.
한계 (Caveats)
- **Dag 수준 재시도 사이의
UP_FOR_RETRY대기. ** 실패한 태스크는 다음 clear를 트리거하기 전에UP_FOR_RETRY에서retry_delay(기본 5분) 동안 대기한 뒤FAILED로 전환되므로, 각 Dag 수준 재시도는 최소 그 간격만큼 지연돼요. - **멱등성은 여러분의 책임이에요. ** 지워진 모든 태스크가 다시 실행돼요. 실패한 시도에서 외부 부작용을 만들었다면 — 행(row) 쓰기, 메시지 전송, 파일 업로드 — 태스크가 멱등하지 않다면(예: upsert, 중복 제거 키, 조건부 쓰기 사용) 다시 생성돼요.
- **루프 위험. ** 카운터나 실패하는 태스크의 버그가 무한 clear-and-retry 루프를 일으킬 수 있어요. 항상 시도 횟수에 상한을 두고 상한에 도달하면 알림을 받으세요.
- **비용. ** Dag 수준 재시도는 지워진 모든 태스크를 처음부터 다시 실행해요. 시간과 리소스 비용이 허용 가능한지 확인하고, 시도 상한을 낮게 유지하세요.
- **실행 중인 태스크. ** Dag에 병렬 분기가 있다면, 한 태스크가 실패하는 동안 같은 Dag run의 다른 태스크가 아직 실행 중일 수 있어요. 다른 작업이 진행 중일 때 run을 지우는 것이 안전하도록 Dag와 그 부작용을 설계하세요.
- **토큰 수명. ** 대부분의 auth manager가 발급한 토큰은 만료돼요(예: simple auth manager의 JWT는 기본 24시간). Variable에 저장된 정적 토큰은 만료 후 조용히 401을 반환해요. 루프 가드가 잘못된 동작을 막지만 Dag 수준 재시도는 동작을 멈춰요. 장기 실행 배포에서는 자격 증명을 주기적으로 갱신하거나 더 긴 수명의 서비스 자격 증명을 사용하세요.