Airflow 해부 (7) 응용: 재학습 파이프라인 DAG
집계, 학습, 평가, 승격 판단으로 이어지는 주간 재학습 DAG를 작성합니다. MLflow의 승격 자동화 함수를 마지막 task로 붙이고, 시리즈 전체의 개념이 한 DAG에서 어떻게 만나는지 정리하며 마칩니다.
Airflow 해부 (7) 응용: 재학습 파이프라인 DAG
Airflow 해부 시리즈의 마지막, 7편입니다. 전체 목차는 0편에 있습니다.
만들 것
매주 월요일 새벽에 도는 재학습 파이프라인입니다. 지난 한 주 데이터를 집계하고, 모델을 학습하고, 검증셋으로 평가한 뒤, MLflow 해부 5편에서 만든 promote_if_better로 champion보다 좋을 때만 승격합니다.
1
aggregate ──> train ──> evaluate ──> promote_or_skip
DAG 전체
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
import pendulum
from datetime import timedelta
from airflow.decorators import dag, task
@dag(
schedule="0 3 * * 1", # 매주 월요일 KST 03:00
start_date=pendulum.datetime(2026, 7, 1, tz="Asia/Seoul"),
catchup=False,
default_args={"retries": 2, "retry_delay": timedelta(minutes=10)},
)
def weekly_retrain():
@task
def aggregate(ds: str) -> str:
from src.features import build_weekly_features
return build_weekly_features(end_date=ds) # 경로 반환
@task
def train(feature_path: str, ds: str) -> str:
import mlflow
from src.train import fit_model
mlflow.set_experiment("demand-weekly")
with mlflow.start_run(run_name=f"weekly-{ds}") as run:
fit_model(feature_path) # autolog가 기록
mlflow.set_tag("trigger", "airflow")
return run.info.run_id
@task
def evaluate(run_id: str, feature_path: str) -> str:
from src.evaluate import log_validation_metrics
log_validation_metrics(run_id, feature_path) # 검증 지표를 run에 추가
return run_id
@task
def promote_or_skip(run_id: str) -> bool:
from src.registry import promote_if_better
return promote_if_better(run_id)
ds = "{{ ds }}"
feature_path = aggregate(ds=ds)
run_id = train(feature_path, ds=ds)
promote_or_skip(evaluate(run_id, feature_path))
weekly_retrain()
설계가 시리즈 개념 위에 서 있다
이 짧은 DAG에 앞 편들의 개념이 전부 들어 있습니다.
- logical date (4편): 집계 기준일을
ds로 받습니다. 지난주 파이프라인에 문제가 있었으면 backfill로 그 주만 다시 돌릴 수 있습니다 - XCom은 경로와 id만 (3편): task 사이를 이동하는 것은 feature 경로와 MLflow run_id뿐입니다. 데이터와 모델은 스토리지와 MLflow에 있습니다
- task 안에서 import (1편, 6편):
src.train같은 무거운 import를 task 함수 안에 둬서 scheduler 파싱을 가볍게 유지합니다 - 멱등성 (6편): aggregate는 해당 주 파티션을 덮어쓰고, train은 재실행 시 새 run을 만들 뿐이며, 승격 판단은 지표 비교라 몇 번을 돌려도 안전합니다
- 재시도 (3편): 일시적 장애는 10분 간격 2회 재시도가 흡수합니다
MLflow 쪽에서 보면, Airflow는 MLflow 해부 5편 승격 함수의 “누가 언제 부르는가”를 채워주는 존재입니다. 두 도구의 결합부가 run_id 하나라는 점이 이 구성의 유지보수를 쉽게 만듭니다.
실제 프로젝트에서는
이 DAG는 골격이고, 실제로 붙일 때 추가되는 것들이 있습니다.
- 데이터 도착을 기다리는 Sensor(5편)가 aggregate 앞에 붙습니다
- 승격 결과와 실패 알림이 Slack 콜백(5편)으로 나갑니다
- 학습이 무거우면 train task만 KubernetesPodOperator(5편)로 GPU 노드에 보냅니다
수집부터 모니터링까지 포함한 전체 구성은 NYC 택시 프로젝트 8편에 있습니다.
시리즈를 마치며
일곱 편의 요약입니다.
- Airflow는 scheduler가 metadata DB를 보고 executor에 task를 넘기는 순환 구조입니다 (1편)
- 설치는 constraint와 함께, 운영은 컨테이너로 합니다 (2편, 6편)
- 의존성은 함수 호출로 선언하고, XCom에는 작은 값만 태웁니다 (3편)
- 날짜는 logical date 기준으로 짜고, catchup=False를 기본으로, 과거 재처리는 backfill로 합니다 (4편)
- 접속 정보는 Connection으로 코드에서 분리합니다 (5편)
- 멱등성이 재시도와 backfill을 안전하게 만듭니다 (6편)
같은 글감의 실험 추적 도구 편은 MLflow 해부 시리즈입니다.
이 기사는 저작권자의 CC BY 4.0 라이센스를 따릅니다.