5.1 오케스트레이션이란?
오케스트레이션은 여러 데이터 파이프라인 작업을 조율하고 스케줄링하는 것입니다. 마치 오케스트라 지휘자가 각 악기의 연주를 조율하듯이, 오케스트레이터는 각 데이터 작업의 실행을 관리합니다.
5.2 Apache Airflow
가장 인기있는 오픈소스 워크플로우 관리 플랫폼:
DAG (Directed Acyclic Graph)
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'data-team',
'retries': 3,
'retry_delay': timedelta(minutes=5)
}
with DAG(
'daily_sales_pipeline',
default_args=default_args,
schedule_interval='0 2 * * *', # 매일 새벽 2시
start_date=datetime(2025, 1, 1),
catchup=False
) as dag:
extract = PythonOperator(
task_id='extract_orders',
python_callable=extract_orders
)
transform = PythonOperator(
task_id='transform_orders',
python_callable=transform_orders
)
load = PythonOperator(
task_id='load_to_warehouse',
python_callable=load_to_warehouse
)
notify = PythonOperator(
task_id='send_notification',
python_callable=send_slack_notification
)
# 의존성 정의
extract >> transform >> load >> notify
센서 (Sensor)
외부 이벤트를 기다림:
from airflow.sensors.s3 import S3KeySensor
wait_for_file = S3KeySensor(
task_id='wait_for_data_file',
bucket_name='my-bucket',
bucket_key='data/{{ ds }}/file.csv',
poke_interval=60, # 60초마다 체크
timeout=3600 # 1시간 타임아웃
)
동적 DAG
tables = ['users', 'orders', 'products']
for table in tables:
extract_task = PythonOperator(
task_id=f'extract_{table}',
python_callable=extract_table,
op_kwargs={'table_name': table}
)
transform_task = PythonOperator(
task_id=f'transform_{table}',
python_callable=transform_table,
op_kwargs={'table_name': table}
)
extract_task >> transform_task
5.3 다른 오케스트레이터
Prefect
from prefect import flow, task
@task
def extract():
return fetch_data()
@task
def transform(data):
return clean_data(data)
@task
def load(data):
save_to_warehouse(data)
@flow
def etl_pipeline():
data = extract()
cleaned = transform(data)
load(cleaned)
Dagster
from dagster import asset
@asset
def orders_raw():
return fetch_orders()
@asset
def orders_clean(orders_raw):
return clean_orders(orders_raw)
@asset
def daily_sales(orders_clean):
return aggregate_sales(orders_clean)
5.4 스케줄링 전략
시간 기반
- 매일:
0 2 * * * - 매주 월요일:
0 2 * * 1 - 매월 1일:
0 2 1 * *
이벤트 기반
- 파일 도착 시
- API 호출 시
- 다른 파이프라인 완료 시
5.5 에러 핸들링
재시도
task = PythonOperator(
task_id='flaky_task',
python_callable=unreliable_function,
retries=3,
retry_delay=timedelta(minutes=5),
retry_exponential_backoff=True
)
알림
def on_failure(context):
send_slack_alert(f"Task {context['task_instance']} failed!")
task = PythonOperator(
task_id='important_task',
python_callable=critical_function,
on_failure_callback=on_failure
)
5.6 의존성 관리
선형 의존성
extract >> transform >> load
병렬 실행
extract >> [transform_a, transform_b, transform_c] >> combine >> load
조건부 실행
from airflow.operators.python import BranchPythonOperator
def decide_branch(**context):
if data_ready():
return 'process_data'
else:
return 'wait_task'
branch = BranchPythonOperator(
task_id='decide',
python_callable=decide_branch
)
5.7 오케스트레이터 비교
| 특징 | Airflow | Prefect | Dagster |
|---|---|---|---|
| 출시년도 | 2014 (Airbnb) | 2018 | 2019 |
| 핵심 개념 | DAG, Operator | Flow, Task | Asset, Op |
| UI | 강력함 | 현대적 | 데이터 중심 |
| 학습 곡선 | 가파름 | 완만함 | 중간 |
| 클라우드 지원 | AWS MWAA | Prefect Cloud | Dagster Cloud |
| 동적 워크플로우 | 제한적 | 우수 | 우수 |
| 데이터 품질 | 별도 구현 | 별도 구현 | 내장 |
| 커뮤니티 | 매우 큼 | 성장 중 | 성장 중 |
언제 무엇을 선택할까?
Airflow를 선택하세요:
- 검증된 솔루션 필요
- 대규모 조직
- 복잡한 의존성 그래프
- 많은 커뮤니티 플러그인 활용
- AWS MWAA 사용 예정
Prefect를 선택하세요:
- 빠른 개발 속도 중요
- 동적 워크플로우 필요
- Python 네이티브 경험 선호
- Negative engineering (실패 처리 자동화)
Dagster를 선택하세요:
- 데이터 품질 중요
- 자산 중심 사고 (Asset-centric)
- 개발-프로덕션 일관성
- 타입 안전성 중요
5.8 고급 DAG 패턴
팬아웃/팬인 패턴
from airflow import DAG
from airflow.operators.python import PythonOperator
with DAG('fan_out_fan_in') as dag:
start = PythonOperator(task_id='start', python_callable=start_job)
# Fan-out: 병렬 처리
process_tasks = []
for i in range(10):
task = PythonOperator(
task_id=f'process_{i}',
python_callable=process_partition,
op_kwargs={'partition': i}
)
process_tasks.append(task)
# Fan-in: 결과 수집
combine = PythonOperator(task_id='combine', python_callable=combine_results)
# 의존성
start >> process_tasks >> combine
서브DAG 패턴
from airflow import DAG
from airflow.operators.subdag import SubDagOperator
def create_processing_subdag(parent_dag_id, child_dag_id, args):
with DAG(f'{parent_dag_id}.{child_dag_id}', default_args=args) as dag:
extract = PythonOperator(task_id='extract', python_callable=extract_data)
transform = PythonOperator(task_id='transform', python_callable=transform_data)
load = PythonOperator(task_id='load', python_callable=load_data)
extract >> transform >> load
return dag
with DAG('main_pipeline') as dag:
start = PythonOperator(task_id='start', python_callable=start)
process_users = SubDagOperator(
task_id='process_users',
subdag=create_processing_subdag('main_pipeline', 'process_users', args)
)
process_orders = SubDagOperator(
task_id='process_orders',
subdag=create_processing_subdag('main_pipeline', 'process_orders', args)
)
start >> [process_users, process_orders]
TaskGroup 패턴 (Airflow 2.0+)
from airflow.utils.task_group import TaskGroup
with DAG('task_group_example') as dag:
start = PythonOperator(task_id='start', python_callable=start)
with TaskGroup('processing_group') as processing:
extract = PythonOperator(task_id='extract', python_callable=extract)
transform = PythonOperator(task_id='transform', python_callable=transform)
load = PythonOperator(task_id='load', python_callable=load)
extract >> transform >> load
with TaskGroup('validation_group') as validation:
check_schema = PythonOperator(task_id='check_schema', python_callable=check_schema)
check_counts = PythonOperator(task_id='check_counts', python_callable=check_counts)
end = PythonOperator(task_id='end', python_callable=end)
start >> processing >> validation >> end
5.9 스케줄링 고급 기법
Cron 표현식 완벽 가이드
| 패턴 | 의미 | 예시 |
|---|---|---|
| 0 2 * * * | 매일 새벽 2시 | 일일 배치 |
| */15 * * * * | 15분마다 | 실시간 동기화 |
| 0 */4 * * * | 4시간마다 | 주기적 수집 |
| 0 9 * * 1 | 매주 월요일 9시 | 주간 리포트 |
| 0 0 1 * * | 매월 1일 자정 | 월간 집계 |
| 0 0 1 1 * | 매년 1월 1일 | 연간 보고서 |
| 0 9-17 * * 1-5 | 평일 9시-17시 매시 | 업무시간 모니터링 |
Timetable (Airflow 2.2+)
from airflow.timetables.base import Timetable
from pendulum import DateTime
class CustomBusinessDayTimetable(Timetable):
"""평일만 실행하는 커스텀 타임테이블"""
def next_dagrun_info(self, last_automated_data_interval, restriction):
if last_automated_data_interval is None:
# 첫 실행
next_start = restriction.earliest
else:
next_start = last_automated_data_interval.end
# 주말 건너뛰기
while next_start.weekday() >= 5: # 5=토, 6=일
next_start = next_start.add(days=1)
return DagRunInfo(
data_interval=DataInterval(next_start, next_start.add(days=1))
)
# DAG에서 사용
dag = DAG(
'business_day_pipeline',
timetable=CustomBusinessDayTimetable(),
start_date=datetime(2025, 1, 1)
)
데이터 기반 스케줄링
from airflow.sensors.external_task import ExternalTaskSensor
with DAG('downstream_pipeline') as dag:
# 상위 파이프라인 완료 대기
wait_for_upstream = ExternalTaskSensor(
task_id='wait_for_data_extraction',
external_dag_id='upstream_pipeline',
external_task_id='extract_complete',
timeout=3600,
poke_interval=60
)
process = PythonOperator(
task_id='process_data',
python_callable=process_data
)
wait_for_upstream >> process
5.10 파라미터와 변수 관리
Airflow Variables
from airflow.models import Variable
# 변수 설정 (UI 또는 CLI)
# airflow variables set db_connection "postgresql://localhost/db"
# DAG에서 사용
db_conn = Variable.get("db_connection")
api_key = Variable.get("api_key", default_var="default_key")
# JSON 변수
config = Variable.get("pipeline_config", deserialize_json=True)
batch_size = config['batch_size']
DAG Run Parameters
from airflow.operators.python import PythonOperator
def process_with_params(**context):
# DAG run 파라미터 접근
params = context['dag_run'].conf
start_date = params.get('start_date', '2025-01-01')
end_date = params.get('end_date', '2025-12-31')
print(f"Processing {start_date} to {end_date}")
process_task = PythonOperator(
task_id='process',
python_callable=process_with_params,
provide_context=True
)
# Trigger with parameters
# airflow dags trigger my_dag -c '{"start_date": "2025-12-01", "end_date": "2025-12-26"}'
템플릿 변수
from airflow.operators.bash import BashOperator
# Jinja 템플릿 사용
process = BashOperator(
task_id='process_daily',
bash_command="""
python process.py \
--date {{ ds }} \
--year {{ macros.ds_format(ds, '%Y-%m-%d', '%Y') }} \
--month {{ macros.ds_format(ds, '%Y-%m-%d', '%m') }} \
--execution-time {{ ts }}
"""
)
# 일반적인 템플릿 변수:
# {{ ds }} # 2025-12-26
# {{ ds_nodash }} # 20251226
# {{ ts }} # 2025-12-26T10:00:00+00:00
# {{ execution_date }} # DateTime 객체
# {{ prev_ds }} # 이전 실행 날짜
# {{ next_ds }} # 다음 실행 날짜
5.11 에러 복구 전략
재시도 정책
from datetime import timedelta
task = PythonOperator(
task_id='flaky_task',
python_callable=unreliable_function,
# 재시도 설정
retries=5,
retry_delay=timedelta(minutes=5),
retry_exponential_backoff=True,
max_retry_delay=timedelta(hours=1),
# 특정 예외만 재시도
retry_exponential_backoff=True
)
SLA (Service Level Agreement)
def sla_miss_callback(dag, task_list, blocking_task_list, slas, blocking_tis):
"""SLA 위반 시 호출"""
message = f"SLA missed for tasks: {[task.task_id for task in task_list]}"
send_alert(message)
with DAG(
'sla_monitored_pipeline',
sla_miss_callback=sla_miss_callback,
default_args={'sla': timedelta(hours=2)}
) as dag:
task = PythonOperator(
task_id='time_sensitive_task',
python_callable=important_function,
sla=timedelta(hours=1) # 이 태스크는 1시간 SLA
)
Circuit Breaker 패턴
class CircuitBreaker:
def __init__(self, failure_threshold=5, timeout=60):
self.failure_count = 0
self.failure_threshold = failure_threshold
self.timeout = timeout
self.last_failure_time = None
self.state = 'CLOSED' # CLOSED, OPEN, HALF_OPEN
def call(self, func, *args, **kwargs):
if self.state == 'OPEN':
if time.time() - self.last_failure_time > self.timeout:
self.state = 'HALF_OPEN'
else:
raise Exception("Circuit breaker is OPEN")
try:
result = func(*args, **kwargs)
if self.state == 'HALF_OPEN':
self.state = 'CLOSED'
self.failure_count = 0
return result
except Exception as e:
self.failure_count += 1
self.last_failure_time = time.time()
if self.failure_count >= self.failure_threshold:
self.state = 'OPEN'
raise e
# 사용
breaker = CircuitBreaker()
def call_external_api():
return breaker.call(requests.get, 'https://api.example.com/data')
5.12 리소스 관리
Pools (동시 실행 제한)
# Pool 생성 (UI 또는 CLI)
# airflow pools set database_pool 5 "Database connection pool"
with DAG('resource_managed_pipeline') as dag:
# 최대 5개 태스크만 동시 실행
task1 = PythonOperator(
task_id='db_task_1',
python_callable=query_database,
pool='database_pool'
)
task2 = PythonOperator(
task_id='db_task_2',
python_callable=query_database,
pool='database_pool',
priority_weight=10 # 높은 우선순위
)
Executor 비교
| Executor | 용도 | 확장성 | 복잡도 |
|---|---|---|---|
| SequentialExecutor | 개발/테스트 | 없음 (1개씩) | 매우 낮음 |
| LocalExecutor | 소규모 프로덕션 | 단일 노드 | 낮음 |
| CeleryExecutor | 대규모 프로덕션 | 다중 워커 | 중간 |
| KubernetesExecutor | 클라우드 네이티브 | 자동 확장 | 높음 |
| DaskExecutor | 데이터 과학 | 동적 확장 | 중간 |
Kubernetes Pod Operator
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator
process_task = KubernetesPodOperator(
task_id='process_in_pod',
name='data-processing-pod',
namespace='data-pipelines',
image='my-company/data-processor:latest',
cmds=['python', 'process.py'],
arguments=['--date', '{{ ds }}'],
env_vars={
'DATABASE_URL': Variable.get('database_url'),
'API_KEY': Variable.get('api_key')
},
resources={
'request_memory': '4G',
'request_cpu': '2',
'limit_memory': '8G',
'limit_cpu': '4'
},
get_logs=True,
is_delete_operator_pod=True
)
5.13 테스팅
DAG 검증 테스트
import pytest
from airflow.models import DagBag
def test_no_import_errors():
"""모든 DAG가 임포트 에러 없이 로드되는지 확인"""
dag_bag = DagBag()
assert len(dag_bag.import_errors) == 0, f"DAG import errors: {dag_bag.import_errors}"
def test_dag_structure():
"""DAG 구조 검증"""
dag_bag = DagBag()
dag = dag_bag.get_dag('my_pipeline')
# 태스크 존재 확인
assert 'extract' in dag.task_dict
assert 'transform' in dag.task_dict
assert 'load' in dag.task_dict
# 의존성 확인
extract = dag.get_task('extract')
assert len(extract.downstream_task_ids) == 1
assert 'transform' in extract.downstream_task_ids
태스크 단위 테스트
from airflow.models import TaskInstance
from airflow.utils.state import State
def test_extract_task():
"""Extract 태스크 단위 테스트"""
# DAG 로드
dag_bag = DagBag()
dag = dag_bag.get_dag('my_pipeline')
# 태스크 실행
task = dag.get_task('extract')
ti = TaskInstance(task=task, execution_date=datetime(2025, 12, 26))
# 실행 및 검증
ti.run(ignore_ti_state=True)
assert ti.state == State.SUCCESS
# XCom에서 결과 확인
result = ti.xcom_pull(task_ids='extract')
assert result is not None
assert len(result) > 0
요약
이 장에서는 데이터 파이프라인 오케스트레이션의 핵심 개념과 실전 기법을 배웠습니다:
- 오케스트레이터 비교: Airflow, Prefect, Dagster의 차이점과 선택 기준
- DAG 설계: 팬아웃/팬인, 서브DAG, TaskGroup 패턴
- 스케줄링: Cron, 커스텀 타임테이블, 데이터 기반 트리거
- 의존성 관리: 선형, 병렬, 조건부 실행
- 에러 핸들링: 재시도, SLA, Circuit Breaker
- 리소스 관리: Pools, Executors, Kubernetes
- 테스팅: DAG 검증, 태스크 단위 테스트
오케스트레이션 베스트 프랙티스
- 멱등성: 재실행 가능하게
- 작은 태스크: 단일 책임 원칙
- 명확한 이름: 무엇을 하는지 명확히
- 모니터링: 실패 알림 설정
- 문서화: DAG 설명 추가
- 파라미터화: 하드코딩 금지
- 테스트: 프로덕션 배포 전 검증
- 리소스 제한: Pool로 동시 실행 관리
복습 문제
- Airflow, Prefect, Dagster의 주요 차이점은 무엇인가요?
- DAG에서 팬아웃/팬인 패턴이 필요한 경우는 언제인가요?
- Cron 표현식 "0 9-17 * * 1-5"의 의미를 설명하세요.
- ExternalTaskSensor가 필요한 이유는 무엇인가요?
- 재시도 정책에서 Exponential Backoff가 중요한 이유는?
- Pool을 사용하는 이유와 활용 사례를 설명하세요.
- KubernetesExecutor와 CeleryExecutor의 차이점은?
- DAG 테스트에서 검증해야 할 항목은 무엇인가요?
실습 과제
- Airflow DAG를 작성하여 ETL 파이프라인을 구현하세요 (Extract, Transform, Load 3단계).
- 동적 DAG를 구현하여 여러 테이블을 병렬로 처리하세요.
- 재시도 정책과 알림을 포함한 프로덕션급 DAG를 작성하세요.
- DAG 검증 테스트를 작성하고 pytest로 실행하세요.