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을 만들 수 있어요. Callback에 kwargs가 지정되면 콜백 함수에 전달돼요. 비동기 콜백은 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과 데드라인에 대한 정보를 포함한
contextkwarg를 콜백에 자동으로 제공해요. 이것을 받으려면 콜백에서**kwargs를 받고kwargs["context"]에 접근하거나,context라는 이름의 파라미터를 추가하면 돼요. context가 필요 없는 콜백은 생략할 수 있어요 — Airflow는 callable이 받는 kwargs만 전달할 거예요.context키워드는 예약되어 있고Callback의kwargs파라미터에 사용할 수 없어요. 시도하면 DAG 파싱 시간에ValueError가 발생해요.
커스텀 동기 콜백은 이렇게 생길 수 있어요:
- 이 메서드를 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
- 이것을 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",
)
커스텀 비동기 콜백은 이렇게 생길 수 있어요:
- 이 메서드를 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
- Triggerer를 재시작하세요.
- 이것을 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_id와run_id만 사용할 수 있어요. 다른 것을 선언하면 데드라인 평가 시ValueError가 발생해요. 참조 자체를 구성하려면 생성자 필드를 주거나 Airflow Variable에서 읽으세요. - 데이터베이스 접근 (Database Access) — 필요한 경우 Airflow 데이터베이스 쿼리에
session파라미터를 사용하세요.