728x90
반응형
SMALL
실습 소스 : 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. Apache Airflow 핵심 개념 정리
- DAG (Directed Acyclic Graph): 방향성 비순환 그래프로, 파이프라인의 작업 순서와 의존성을 정의합니다.
- Operator: 실행할 개별 작업의 주체 (예: Bash 셸 실행, Python 함수 호출, SQL 실행 등).
- Task: Operator가 DAG 내에서 실제 정의된 상태 (노드).
- TaskInstance: 특정 시간/스케줄에 할당되어 실행된 Task의 실제 메모리 객체.
2. BashOperator 기초 (01_basics_bash.py)
dags/01_basics_bash.py 파일은 기본 DAG 작성 문법과 BashOperator 사용법을 보여줍니다.
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
# 1. default_args 공통 옵션 정의
default_args = {
"owner": "aic-de7-admin",
"depends_on_past": False,
"retries": 1,
"retry_delay": timedelta(minutes=5)
}
# 2. DAG 생성 세션
with DAG(
dag_id="01_basics_bash",
description="Airflow DAG 작성 기본형",
default_args=default_args,
schedule_interval="@daily",
start_date=datetime(2026, 6, 29),
catchup=False, # 과거 미실행 구간 소급 실행 금지
tags=['bash', 'basic']
) as dag:
# 3. Task (Operator) 정의
t1 = BashOperator(task_id="date-print", bash_command="date")
t2 = BashOperator(task_id="sleep", bash_command="sleep 3")
t3 = BashOperator(task_id="echo-print", bash_command='echo "hello airflow task"')
# 4. 의존성 (수행 순서) 설정
t1 >> t2 >> t3
💡 핵심 포인트
catchup=False: 과거start_date부터 현재까지 누락된 실행을 소급하여 대량 실행(Backfill)하는 현상을 방지합니다.t1 >> t2 >> t3: 비트 시프트 연산자(>>)로 순차적 실행 흐름을 직관적으로 설정합니다.
3. PythonOperator 및 XCom 데이터 공유 (02_basics_python.py)
Python 함수를 실행하고, Task 간 메시지/작은 데이터를 주고받을 때는 XCom (Cross Communication)을 이용합니다.
dags/02_basics_python.py에서 검증할 수 있습니다.
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import timedelta
import logging
import pendulum
KST = pendulum.timezone("Asia/Seoul")
# Extract 콜백 함수
def _extract_cb(**kwargs):
ti = kwargs['ti'] # TaskInstance 객체
ds = kwargs['ds'] # YYYY-MM-DD
run_id = kwargs['run_id']
logging.info(f"Extract 작업 수행: ds={ds}, run_id={run_id}")
# return 문으로 반환된 데이터는 XCom에 자동으로 저장(push)됩니다.
return f"ds = {ds}, run_id = {run_id}"
# Transform 콜백 함수
def _transform_cb(**kwargs):
ti = kwargs["ti"]
# extract_task가 XCom으로 반환한 데이터 가져오기(pull)
data = ti.xcom_pull(task_ids="extract_task")
logging.info(f"Transform에서 전달받은 XCom 데이터: {data}")
with DAG(
dag_id="02_basics_python",
default_args={"owner": "aic-de7-admin", "retries": 1},
schedule_interval="@once",
start_date=pendulum.datetime(2026, 6, 29, tz=KST),
catchup=False,
tags=['python', 'xcom']
) as dag:
extract_task = PythonOperator(
task_id="extract_task",
python_callable=_extract_cb
)
transform_task = PythonOperator(
task_id="transform_task",
python_callable=_transform_cb
)
extract_task >> transform_task
💡 XCom 작동 원리
PythonOperator콜백 함수에서return한 값은 기본적으로key="return_value"로 XCom 공간에 자동 등록됩니다.- 후속 Task는
ti.xcom_pull(task_ids='대상_task_id')를 호출하여 이전 Task의 가공 결과를 받아올 수 있습니다.
4. 2편 요약 및 다음 편 예고
이번 글에서는 Airflow의 DAG 구성 요소, BashOperator, PythonOperator, 그리고 XCom 데이터 교환을 실습해 보았습니다.
다음 3편에서는 Airflow의 생산성을 획기적으로 올려주는 Jinja 템플릿, 크론 스케줄링, 그리고 BranchPythonOperator를 이용한 조건 분기를 알아보겠습니다.
동작방법
동작하는 것을 보고 싶으면
브라우저 키셔서 localhost:8000 접속하셔서
id/pwd : airflow/airflow
로그인하시고 DAGs 에서 저희가 만든 DAG 들어가서 실행해보면 됩니다.

반응형
LIST