Deadline Alerts

Deadline Alerts

Dag run에 시간 임계값을 설정하고 그 값을 초과했을 때 자동으로 대응하는 Deadline Alerts 기능을 설명하는 문서예요. 참조점(reference), 간격(interval), 콜백(callback) 세 가지 필수 파라미터로 Deadline Alert를 만들고, 내장 참조와 콜백, 커스텀 참조와 콜백, 여러 Deadline Alert 구성까지 코드와 함께 살펴볼게요.

출처: 문서

본문

경고 (Warning)

Deadline Alerts는 Airflow 3.1의 새로운 기능으로 실험적(experimental)으로 간주해야 해요. 이 기능은 사용자 피드백에 따라 향후 버전에서 경고 없이 변경될 수 있어요.

이것은 실험적 기능이에요.

Deadline Alerts를 사용하면 Dag run에 시간 임계값을 설정하고 그 임계값을 초과했을 때 자동으로 대응할 수 있어요. 참조점을 선택하고 간격을 설정하고 데드라인을 놓치면 실행할 콜백을 정의해 Deadline Alerts를 구성해요. 참조는 dagrun이 큐에 들어갈 때 같은 내장 DeadlineReference 옵션 중 하나이거나 타임스탬프를 반환하는 어떤 커스텀 메서드일 수 있어요. 콜백은 Airflow의 Notifier 중 하나이거나 커스텀 콜백 함수일 수 있어요.

SLA에서 마이그레이션하기 (Migrating from SLA)

SLA에서 Deadlines로 마이그레이션하는 데 도움이 필요하면 마이그레이션 가이드를 참고하세요.

Deadline Alert 만들기 (Creating a Deadline Alert)

Deadline Alert를 만들려면 세 가지 필수 파라미터가 필요해요:

  • Reference: 언제부터 카운트를 시작할지
  • Interval: 참조점에서 얼마나 앞/뒤에서 알림을 트리거할지(timedelta 또는 VariableInterval 같은 동적 간격)
  • Callback: 콜백 객체로, callable의 경로와 데드라인 초과 시 전달할 선택적 kwargs를 포함

Deadlines가 계산되는 방식은 다음과 같아요:

[Reference] ------ [Interval] ------> [Deadline]
    ^                                     ^
    |                                     |
 Start time                          Trigger point

아래는 예시 Dag 구현이에요. Dag가 큐에 들어간 지 15분 안에 끝나지 않으면 Slack 메시지를 보내요:

from datetime import datetime, timedelta
from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference
from airflow.providers.slack.notifications.slack_webhook import SlackWebhookNotifier
from airflow.providers.standard.operators.empty import EmptyOperator

with DAG(
    dag_id="deadline_alert_example",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_QUEUED_AT,
        interval=timedelta(minutes=15),
        callback=AsyncCallback(
            SlackWebhookNotifier,
            kwargs={
                "text": "🚨 Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}"
            },
        ),
    ),
):
    EmptyOperator(task_id="example_task")

이 예시의 타임라인은 다음과 같아요:

|------|-----------|---------|-----------|--------|
    Scheduled    Queued    Started    Deadline
     00:00       00:03      00:05      00:18

참고 (Note)

AsyncCallback의 import 경로가 Airflow 3.2에서 airflow.sdk.definitions.deadline에서 airflow.sdk로 변경됐어요.

내장 참조 사용하기 (Using Built-in References)

Airflow는 DeadlineAlert와 함께 사용할 수 있는 몇 가지 내장 참조점을 제공해요.

DeadlineReference.DAGRUN_QUEUED_AT

Dag run이 큐에 들어간 시점부터 시간을 측정해요. 리소스 제약을 모니터링하는 데 유용해요.

DeadlineReference.DAGRUN_LOGICAL_DATE

Dag run이 실행되도록 스케줄된 시점을 참조해요. 예를 들어 timedelta(minutes=15) 간격을 설정하면, Dag가 실제로 실행을 시작했는지(또는 언제 시작했는지)와 무관하게 스케줄된 시작 15분 후에 완료되지 않았으면 알림이 트리거돼요. 스케줄된 Dags가 다음 스케줄 실행 전에 완료되도록 보장하는 데 유용해요.

DeadlineReference.FIXED_DATETIME

시간상의 고정 지점을 지정해요. Dags가 특정 시간까지 완료되어야 할 때 유용해요.

DeadlineReference.AVERAGE_RUNTIME

이전 성공한 Dag run의 평균 실행 시간을 기반으로 데드라인을 계산해요. 이 참조는 과거 실행 데이터를 분석해 현재 run이 언제 완료되어야 하는지 예측해요. 데드라인은 현재 시간 + 계산된 평균 실행 시간 + 간격으로 설정돼요. 과거 데이터가 충분하지 않으면 데드라인이 생성되지 않아요.

파라미터:

  • max_runs (int, optional): 분석할 최근 성공한 Dag run의 최대 수. 기본값 10.
  • min_runs (int, optional): 평균을 계산하는 데 필요한 최근 성공한 Dag run의 최소 수. 기본값은 max_runs와 같음.

사용 예시:

# Use default settings (analyze up to 10 runs, require 10 runs)
DeadlineReference.AVERAGE_RUNTIME()

# Analyze up to 20 runs but calculate with minimum 5 runs
DeadlineReference.AVERAGE_RUNTIME(max_runs=20, min_runs=5)

# Strict: require exactly 15 runs to calculate
DeadlineReference.AVERAGE_RUNTIME(max_runs=15, min_runs=15)

평균 실행 시간을 사용하는 예시:

with DAG(
    dag_id="average_runtime_deadline",
    deadline=DeadlineAlert(
        reference=DeadlineReference.AVERAGE_RUNTIME(max_runs=15, min_runs=5),
        interval=timedelta(minutes=30),  # Alert if 30 minutes past average runtime
        callback=AsyncCallback(
            SlackWebhookNotifier,
            kwargs={"text": "🚨 Dag {{ dag_run.dag_id }} is running longer than expected!"},
        ),
    ),
):
    EmptyOperator(task_id="data_processing")

계산된 과거 평균이 30분이었다면 이 예시의 타임라인은 다음과 같아요:

|------|----------|--------------|--------------|--------|
     Queued     Start            |           Deadline
     09:00      09:05          09:35          10:05
                  |              |              |
                  |--- Average --|-- Interval --|
                       (30 min)      (30 min)

고정 datetime을 사용하는 예시:

tomorrow_at_ten = datetime.combine(datetime.now().date() + timedelta(days=1), time(10, 0))

with DAG(
    dag_id="fixed_deadline_alert",
    deadline=DeadlineAlert(
        reference=DeadlineReference.FIXED_DATETIME(tomorrow_at_ten),
        interval=timedelta(minutes=-30),  # Alert 30 minutes before the reference.
        callback=AsyncCallback(
            SlackWebhookNotifier,
            kwargs={
                "text": "🚨 Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}"
            },
        ),
    ),
):
    EmptyOperator(task_id="example_task")

이 예시의 타임라인은 다음과 같아요:

|------|----------|---------|------------|--------|
     Queued     Start    Deadline    Reference
     09:15      09:17     09:30       10:00

참고 (Note)

이 경우 간격이 음수 값이므로 데드라인은 참조보다 앞에 있어요.

콜백 사용하기 (Using Callbacks)

데드라인이 초과되면 콜백의 callable이 지정된 kwargs와 함께 실행돼요. 기존 Notifier를 사용하거나 커스텀 callable을 만들 수 있어요. 콜백은 AsyncCallback 또는 SyncCallback(SyncCallback 지원은 3.2에 추가됨)이어야 해요.

내장 Notifier 사용하기 (Using Built-in Notifiers)

Dag run이 큐에 들어간 지 30분 안에 끝나지 않으면 비동기 콜백으로 Slack Notifier를 사용하는 예시예요. 콜백은 Triggerer에서 실행돼요:

with DAG(
    dag_id="slack_deadline_alert_async",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_QUEUED_AT,
        interval=timedelta(minutes=30),
        callback=AsyncCallback(
            SlackWebhookNotifier,
            kwargs={
                "text": "🚨 Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}"
            },
        ),
    ),
):
    EmptyOperator(task_id="example_task")

동기 콜백을 사용하는 같은 예시예요. 콜백은 executor에서 실행돼요:

with DAG(
    dag_id="slack_deadline_alert_sync",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_QUEUED_AT,
        interval=timedelta(minutes=30),
        callback=SyncCallback(
            SlackWebhookNotifier,
            kwargs={
                "text": "🚨 Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}"
            },
        ),
    ),
):
    EmptyOperator(task_id="example_task")

커스텀 콜백 만들기 (Creating Custom Callbacks)

더 복잡한 처리를 위해 커스텀 callable을 만들 수 있어요. Callbackkwargs가 지정되면 콜백 함수에 전달돼요. 비동기 콜백은 Triggerer의 시스템 경로의 어딘가에 정의되어야 해요. 동기 콜백은 실행될 worker에서 import 가능해야 해요.

참고 (Note) — Async 커스텀 Deadline 콜백 관련:

  • Async 콜백은 Triggerer가 실행하므로, 사용자는 콜백이 Triggerer에서 import 가능하도록 보장해야 해요.
  • 이를 위한 쉬운 방법 중 하나는 plugins 폴더의 새 파일에 callable을 최상위 메서드로 두는 것이에요. 중첩 callable은 현재 지원되지 않아요.
  • 콜백이 추가되거나 변경되면 파일을 다시 로드하기 위해 Triggerer를 재시작해야 해요.

참고 (Note) — Synchronous 콜백 관련:

  • Sync 콜백은 executor로 보내져 최고 우선순위의 Dag 태스크처럼 취급돼요.

참고 (Note)Airflow context:

데드라인이 놓치면 Airflow는 Dag run과 데드라인에 대한 정보를 포함한 context kwarg를 콜백에 자동으로 제공해요. 이것을 받으려면 콜백에서 **kwargs를 받고 kwargs["context"]에 접근하거나, context라는 이름의 파라미터를 추가하면 돼요. context가 필요 없는 콜백은 생략할 수 있어요 — Airflow는 callable이 받는 kwargs만 전달할 거예요. context 키워드는 예약되어 있고 Callbackkwargs 파라미터에 사용할 수 없어요. 시도하면 DAG 파싱 시간에 ValueError가 발생해요.

커스텀 동기 콜백은 이렇게 생길 수 있어요:

  1. 이 메서드를 plugins 폴더(예: $AIRFLOW_HOME/plugins/deadline_callbacks.py)에 두세요:
def custom_sync_callback(**kwargs):
    """Handle deadline violation with custom logic."""
    context = kwargs.get("context", {})
    print(f"Deadline exceeded for Dag {context.get('dag_run', {}).get('dag_id')}!")
    print(f"Context: {context}")
    print(f"Alert type: {kwargs.get('alert_type')}")
    # Additional custom handling here
  1. 이것을 Dag 파일에 두세요:
from datetime import timedelta

from deadline_callbacks import custom_sync_callback

from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.sdk import DAG, DeadlineAlert, DeadlineReference, SyncCallback

with DAG(
    dag_id="custom_sync_deadline_alert",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_QUEUED_AT,
        interval=timedelta(minutes=15),
        callback=SyncCallback(
            custom_sync_callback,
            kwargs={"alert_type": "time_exceeded"},
        ),
    ),
):
    EmptyOperator(task_id="example_task")

팁 (Tip)

SyncCallback은 특정 executor를 대상으로 하는 선택적 executor 파라미터를 받아요. 지정하지 않으면 기본 executor가 사용돼요.

SyncCallback(
    my_callback,
    kwargs={"msg": "deadline missed"},
    executor="celery_executor",
)

커스텀 비동기 콜백은 이렇게 생길 수 있어요:

  1. 이 메서드를 plugins 폴더(예: $AIRFLOW_HOME/plugins/deadline_callbacks.py)에 두세요:
async def custom_async_callback(**kwargs):
    """Handle deadline violation with custom logic."""
    context = kwargs.get("context", {})
    print(f"Deadline exceeded for Dag {context.get('dag_run', {}).get('dag_id')}!")
    print(f"Context: {context}")
    print(f"Alert type: {kwargs.get('alert_type')}")
    # Additional custom handling here
  1. Triggerer를 재시작하세요.
  2. 이것을 Dag 파일에 두세요:
from datetime import timedelta

from deadline_callbacks import custom_async_callback

from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference

with DAG(
    dag_id="custom_deadline_alert",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_QUEUED_AT,
        interval=timedelta(minutes=15),
        callback=AsyncCallback(
            custom_async_callback,
            kwargs={"alert_type": "time_exceeded"},
        ),
    ),
):
    EmptyOperator(task_id="example_task")

템플릿과 컨텍스트 (Templating and Context)

현재 비교적 간단한 버전의 Airflow context가 callable에 전달되고, Airflow는 kwargs에 대해 Jinja Templating을 실행하지 않아요. 다만 Notifier는 실행의 일부로 제공된 context로 이미 템플릿을 실행해요. 즉 Notifier를 사용할 때는 템플릿되는 변수가 단순화된 context에 포함되어 있기만 하면 템플릿을 사용할 수 있어요. 여기에는 현재 Deadline Alert의 ID와 계산된 데드라인 시간, 그리고 Dag Run에 대한 GET REST API 응답에 포함된 데이터가 포함돼요. 더 포괄적인 context와 템플릿 지원은 향후 버전에서 추가될 거예요.

데드라인 계산 (Deadline Calculation)

데드라인의 트리거 시간은 reference가 반환한 datetime에 interval을 더해 계산돼요. FIXED_DATETIME 참조의 경우 음수 간격이 참조 시간 이전에 콜백을 트리거하는 데 특히 유용할 수 있어요.

다음 예시에서 notify_team은 elsewhere에서 정의된 SyncCallback 또는 AsyncCallback이에요:

next_meeting = datetime(2025, 6, 26, 9, 30)

DeadlineAlert(
    reference=DeadlineReference.FIXED_DATETIME(next_meeting),
    interval=timedelta(hours=-2),
    callback=notify_team,
)

이것은 다음 회의가 시작되기 2시간 전에 알림을 트리거해요.

DAGRUN_LOGICAL_DATE의 경우 간격은 일반적으로 양수이며, Dag가 실행되도록 스케줄된 시점을 기준으로 데드라인을 설정해요. 예시:

DeadlineAlert(
    reference=DeadlineReference.DAGRUN_LOGICAL_DATE,
    interval=timedelta(hours=1),
    callback=notify_team,
)

이 경우 Dag가 매일 자정에 실행되도록 스케줄되어 있다면, 오전 1시까지 완료되지 않으면 데드라인이 트리거돼요. 이는 스케줄된 job이 의도된 시작 시간 이후 특정 시간대 내에 완료되도록 보장하는 데 유용해요.

다른 참조를 양수·음수 간격과 결합하는 유연성 덕분에 다양한 운영 요구사항에 맞는 데드라인을 만들 수 있어요.

커스텀 참조 (Custom References)

내장 참조가 대부분의 흔한 시나리오를 처리해요. 하지만 캘린더나 다른 데이터 소스 같은 특정 통합을 위해 커스텀 참조를 만들어야 할 수도 있어요. 그러려면 BaseDeadlineReference에서 상속받는 클래스를 만들고 @deadline_reference 데코레이터를 추가하고 _evaluate_with() 메서드를 구현해요.

데코레이터는 괄호 유무와 관계없이 사용할 수 있어요. 괄호 없이(또는 빈 괄호로) 사용하면 새 Dag run이 생성될 때 참조가 평가돼요. 다른 시점을 선택하려면 DeadlineReference.TYPES 값을 전달하세요.

커스텀 참조 만들기 (Creating a Custom Reference)

from sqlalchemy.orm import Session

from airflow.sdk import BaseDeadlineReference, DeadlineReference, deadline_reference
from airflow.sdk.timezone import datetime

# By default, the evaluate_with method will be executed when the dagrun is created.
@deadline_reference()
class MyCustomDecoratedReference(BaseDeadlineReference):
    """A custom reference evaluated when Dag runs are created."""

    def _evaluate_with(self, *, session: Session, **kwargs) -> datetime:
        # Add your business logic here
        return your_datetime

# You can specify when evaluate_with will be called by providing a DeadlineReference.TYPES value.
@deadline_reference(DeadlineReference.TYPES.DAGRUN_QUEUED)
class MyQueuedReference(BaseDeadlineReference):
    """A custom reference evaluated when Dag runs are queued."""

    # Ask for the Dag run context values supplied by Airflow; see notes below.
    required_kwargs = {"dag_id", "run_id"}

    def _evaluate_with(self, *, session: Session, **kwargs) -> datetime:
        dag_id = kwargs["dag_id"]
        run_id = kwargs["run_id"]
        # Use dag_id and run_id in your calculation
        return your_datetime

Dag에서 커스텀 참조 사용하기 (Using a Custom Reference in a Dag)

등록되면 다른 참조와 마찬가지로 Dag 정의에서 커스텀 참조를 사용해요:

from datetime import timedelta
from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference

with DAG(
    dag_id="custom_reference_example",
    deadline=DeadlineAlert(
        reference=DeadlineReference.MyCustomDecoratedReference,
        interval=timedelta(hours=2),
        callback=AsyncCallback(my_callback),
    ),
):
    # Your tasks here
    ...

여러 Deadline Alerts (Multiple Deadline Alerts)

Dag는 여러 Deadline Alerts를 가질 수 있어요. 단일 DeadlineAlert 대신 deadline 파라미터에 목록을 전달해요. 목록의 각 알림은 독립적으로 평가되며, 각각 참조점과 콜백 타입(sync 또는 async)의 어떤 조합도 사용할 수 있어요.

from datetime import timedelta
from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference, SyncCallback
from airflow.providers.slack.notifications.slack_webhook import SlackWebhookNotifier
from airflow.providers.standard.operators.empty import EmptyOperator

with DAG(
    dag_id="multiple_deadline_alerts",
    deadline=[
        # First alert: warn via Slack (async) if not done 30 min after queuing
        DeadlineAlert(
            reference=DeadlineReference.DAGRUN_QUEUED_AT,
            interval=timedelta(minutes=30),
            callback=AsyncCallback(
                SlackWebhookNotifier,
                kwargs={"text": "⚠️ Dag {{ dag_run.dag_id }} is approaching its deadline."},
            ),
        ),
        # Second alert: escalate via custom sync callback if not done 60 min after queuing
        DeadlineAlert(
            reference=DeadlineReference.DAGRUN_QUEUED_AT,
            interval=timedelta(minutes=60),
            callback=SyncCallback(
                "my_plugins.escalation.escalate_to_oncall",
                kwargs={"severity": "high"},
            ),
        ),
    ],
):
    EmptyOperator(task_id="example_task")

이 패턴은 계층적 알림 전략을 만드는 데 유용해요 — 예를 들어 경고 알림에 이어 Dag가 여전히 실행 중이면 더 긴급한 에스컬레이션을 보내는 방식이에요.

중요 사항 (Important Notes):

  • 시간대 인식 (Timezone Awareness) — 항상 시간대 인식(timezone-aware) datetime 객체를 반환하세요.
  • 무인자 생성 (No-argument Construction) — 커스텀 참조는 등록 중에 인스턴스화되므로 인자 없이 구성할 수 있어야 해요. 참조가 파라미터를 받는다면 @dataclass로 데코레이션하고 모든 필드에 기본값을 주세요.
  • 플러그인 배치 (Plugin Placement) — 커스텀 참조를 두는 편리한 장소 중 하나는 plugins 디렉토리예요.
  • API 서버 재시작 (API Server Restart) — 커스텀 참조를 추가하거나 수정한 후 Airflow API Server를 재시작하세요.
  • 필수 파라미터 (Required Parameters)required_kwargs는 Airflow가 _evaluate_with()로 전달해야 할 Dag run context 값을 선언해요. dag_idrun_id만 사용할 수 있어요. 다른 것을 선언하면 데드라인 평가 시 ValueError가 발생해요. 참조 자체를 구성하려면 생성자 필드를 주거나 Airflow Variable에서 읽으세요.
  • 데이터베이스 접근 (Database Access) — 필요한 경우 Airflow 데이터베이스 쿼리에 session 파라미터를 사용하세요.

더 알아보기 (Learn more)