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. Cron 스케줄링 및 Jinja 템플릿 (03_basics_context_jinja.py)
Airflow에서는 스케줄 매개변수로 숏컷 프리셋(@daily, @hourly) 외에 표준 Cron 표현식(분 시 일 월 요일)을 제공합니다. 또한 명령어나 쿼리에 동적 날짜 변수를 파라미터화할 수 있도록 Jinja 템플릿 및 Macros 기능을 지원합니다.
dags/03_basics_context_jinja.py의 주요 코드입니다.
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import timedelta
import pendulum
KST = pendulum.timezone("Asia/Seoul")
with DAG(
dag_id="03_basics_context_jinja",
description="Macro 및 Jinja 템플릿 사용",
default_args={"owner": "aic-de7-admin", "retries": 1},
schedule_interval="0 9 * * *", # 매일 오전 09시 00분 정각 실행
start_date=pendulum.datetime(2026, 6, 29, tz=KST),
catchup=False,
tags=['macro', 'context', 'jinja']
) as dag:
# 1. ds (YYYY-MM-DD), ti 객체 접근 Jinja 템플릿
t1 = BashOperator(
task_id="jinja_used_task",
bash_command="echo 'DAG 수행시간: {{ ds }}, TaskInstance: {{ ti }}'"
)
# 2. macros 내장 함수 활용 (7일 전 날짜 계산, 난수 생성)
t2 = BashOperator(
task_id="jinja_macro_task",
bash_command="echo '일주일 전: {{ macros.ds_add(ds, -7) }}, 난수: {{ macros.random() }}'"
)
💡 주요 내장 Jinja 템플릿 변수
{{ ds }}: 실행 대상 날짜 (YYYY-MM-DD){{ ds_nodash }}: 하이픈 제외 날짜 (YYYYMMDD){{ macros.ds_add(ds, -7) }}: 특정 날짜 기준 일(day) 수치 가감
2. 동적 분기(Branching) 및 Trigger Rule (04_basics_branching.py)
조건에 따라 특정 작업 경로만 선택해 실행하고 나머지 경로를 생략(Skip)할 경우 BranchPythonOperator를 사용합니다.
dags/04_basics_branching.py를 확인해 보겠습니다.
from airflow import DAG
from airflow.operators.python import PythonOperator, BranchPythonOperator
from airflow.operators.empty import EmptyOperator
from airflow.utils.trigger_rule import TriggerRule
import random, pendulum
KST = pendulum.timezone("Asia/Seoul")
# 조건에 따라 실행할 task_id의 문자열을 반환하는 콜백
def _branch_cb(**kwargs):
if random.choice([True, False]):
return "process" # 'process' task 실행
else:
return "skip" # 'skip' task 실행
with DAG(
dag_id="04_basics_branching",
schedule_interval="@daily",
start_date=pendulum.datetime(2026, 6, 29, tz=KST),
catchup=False,
tags=['branch', 'trigger_rule']
) as dag:
task_start = EmptyOperator(task_id="start")
# 분기 선택 노드
task_branch = BranchPythonOperator(
task_id="branch",
python_callable=_branch_cb
)
task_process = PythonOperator(
task_id="process",
python_callable=lambda **kwargs: print("Process 수행")
)
task_skip = EmptyOperator(task_id="skip")
# 최종 수집 노드 (TriggerRule 적용)
task_end = EmptyOperator(
task_id="end",
# 상위 Task 중 최소 1개만 성공하고 실패가 없으면 정상 진행
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS
)
# 워크플로우 구성
task_start >> task_branch
task_branch >> task_process >> task_end
task_branch >> task_skip >> task_end
┌──► process ──┐
start ─┴─► branch ─────┼──► end (NONE_FAILED_MIN_ONE_SUCCESS)
└──► skip ────┘
💡 TriggerRule 개념
기본적으로 Airflow의 Task는 모든 상위 Task가 success 상태여야 실행됩니다 (all_success).
그러나 Branching으로 인해 하나의 경로가 skipped 상태가 되면, 하위 Task도 생략될 위험이 있습니다.TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS를 부여하면 생략된 경로가 존재하더라도 정상적으로 최종 성공 처리됩니다.
3. 3편 요약 및 다음 편 예고
Jinja 템플릿을 활용한 동적 파라미터 바인딩과 BranchPythonOperator를 통한 조건별 워크플로우 분기 방법을 학습하였습니다.
다음 4편에서는 실제 스마트팩토리 온도 센서 데이터를 생성, 정제하고 DB에 적재하는 [실전 ETL 파이프라인]을 단일 DAG 구조로 개발해 보겠습니다.
'코딩 개발 > Date Engineer' 카테고리의 다른 글
| [DE 실습 5편] [고급] Multi-DAG 디커플링 & FastAPI 외부 REST API 연동 (1) | 2026.09.17 |
|---|---|
| [DE 실습 4편] [실전 ETL] 센서 데이터 수집/정제/MySQL 적재 파이프라인 (0) | 2026.09.16 |
| [DE 실습 2편] Apache Airflow 입문: DAG 구조와 기본 Operator (Bash & Python) (0) | 2026.09.14 |
| [DE 실습 1편] Docker Compose 기반 개발 환경 구축 (Airflow + MySQL + FastAPI) (0) | 2026.09.11 |
| [DE 파이프라인 #4] 실시간 스트리밍 파이프라인 (Kafka & Flink) & ETL vs ELT 아키텍처 심화 비교 (0) | 2026.09.10 |