https://github.com/hosose/DE/tree/airflow
GitHub - hosose/DE at airflow
Contribute to hosose/DE development by creating an account on GitHub.
github.com
1. Why Multi-DAG Architecture? (단일 DAG vs Multi-DAG)
단일 DAG에 수집, 정제, 적재 작업을 모두 몰아넣으면 관리하기 쉽지만 다음과 같은 문제가 발생합니다.
- 특정 단계(예: DB 적재) 실패 시 전체 파이프라인을 재실행해야 함
- 데이터 수집 전담 팀과 정제/분석 팀의 소스코드 충돌
- 모듈화 및 재사용 불가능
Multi-DAG 아키텍처는 Extract DAG, Transform DAG, Load DAG를 각각 독립된 프로세스로 떼어내고 이벤트 트리거 방식으로 연동합니다.
[ DAG 1: Extract ] ──TriggerDagRunOperator──► [ DAG 2: Transform ] ──TriggerDagRunOperator──► [ DAG 3: Load ]
2. Multi-DAG 트리거 구현 (TriggerDagRunOperator)
1) Extract DAG (06_multi_dag_1_extract.py)
수집이 완료되면 TriggerDagRunOperator가 06_multi_dag_2_transform DAG를 자동으로 구동하며 payload(conf)를 통해 파일 경로를 넘겨줍니다.
dags/06_multi_dag_1_extract.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
import pendulum
KST = pendulum.timezone("Asia/Seoul")
with DAG(
dag_id="06_multi_dag_1_extract",
schedule_interval="@daily",
start_date=pendulum.datetime(2026, 6, 29, tz=KST),
catchup=False,
tags=['etl', 'extract']
) as dag:
task_extract = PythonOperator(
task_id="extract",
python_callable=_extract
)
# Transform DAG를 연쇄 실행시키는 트리거 노드
task_trigger_transform = TriggerDagRunOperator(
task_id="trigger_transform",
trigger_dag_id="06_multi_dag_2_transform", # 호출할 대상 DAG ID
conf={
"json_path": "{{ task_instance.xcom_pull(task_ids='extract') }}"
},
reset_dag_run=True,
wait_for_completion=False # 비동기 호출 후 완료
)
task_extract >> task_trigger_transform
2) Transform DAG (06_multi_dag_2_transform.py)
전달받은 dag_run.conf에서 파일 경로를 꺼내 처리 후, 다음 Load DAG를 트리거합니다.
dags/06_multi_dag_2_transform.py
def _transform(**kwargs):
# 상위 DAG가 트리거할 때 전달한 conf 객체에서 파일 경로 수신
dag_run = kwargs["dag_run"]
json_path = dag_run.conf.get('json_path')
# 정제 로직 수행 후 CSV 경로 반환...
return csv_path
# 트리거 생성
task_trigger_load = TriggerDagRunOperator(
task_id="trigger_load",
trigger_dag_id="06_multi_dag_3_load",
conf={"csv_path": "{{ task_instance.xcom_pull(task_ids='transform') }}"},
reset_dag_run=True,
wait_for_completion=False
)
3. 외부 마이크로서비스 (FastAPI) API 연동 패턴
Airflow 파이프라인은 내부 DB 적재에 그치지 않고, 외부 AI/ML 예측 REST API 서비스와 연동할 수 있습니다.
2편에서 구축한 FastAPI 신용평가 API(http://ai-api-server:8000/predict)를 Airflow Task 내부에서 호출하는 패턴 예시입니다.
import requests
def _predict_credit_scores(**kwargs):
# 1. DB 또는 파일에서 신용평가 대상 사용자 수집
users_payload = [
{"user_id": "USER_01", "income": 50000, "loan_amt": 10000},
{"user_id": "USER_02", "income": 80000, "loan_amt": 5000}
]
# 2. FastAPI 마이크로서비스 POST 호출 (Docker 네트워크 이름 활용)
url = "http://ai-api-server:8000/predict"
response = requests.post(url, json=users_payload)
if response.status_code == 200:
predictions = response.json()
logging.info(f"AI 신용평가 응답 결과: {predictions}")
return predictions
else:
raise Exception(f"API 호출 실패 status_code={response.status_code}")
4. 연재 마무리 및 총평
총 6편의 연재를 통해 데이터 엔지니어링 생명주기 이론부터 Docker Compose 실습 인프라 구축, Apache Airflow 기본 문법, 단일 ETL 파이프라인, 그리고 Multi-DAG 모듈화 및 REST API 연동까지 전체 생명주기를 다루었습니다.
그동안 연재를 함께해 주셔서 감사합니다! 🚀