Apache Airflow 2026년 가이드: 파이프라인 오케스트레이션, DAG 및 면접 질문 정리
Apache Airflow 3.2의 Task SDK를 활용한 DAG 구축, 에셋 파티션, 네이티브 비동기 태스크, 데이터 엔지니어 면접에서 자주 출제되는 질문을 종합적으로 다루는 튜토리얼입니다.

Apache Airflow 3.2는 3.0 아키텍처 전면 개편 이후 가장 중요한 진화를 이룬 릴리스입니다. PyPI 월간 다운로드 수가 1,000만 건을 넘어섰으며, Airbnb부터 Spotify까지 다양한 기업에서 채택하고 있어 데이터 파이프라인 구축, 스케줄링, 모니터링 분야에서 핵심 도구로서의 위상을 확고히 하고 있습니다. 이 튜토리얼에서는 새로운 Task SDK를 활용한 DAG 작성, 파이프라인 오케스트레이션 패턴, 그리고 데이터 엔지니어 직무 면접에서 자주 출제되는 질문을 살펴봅니다.
Airflow 3.2는 2026년 4월에 릴리스되었으며, 세분화된 데이터 인식 스케줄링을 위한 에셋 파티션, PythonOperator의 네이티브 비동기 지원, 멀티팀 배포 기능이 도입되었습니다. 모든 DAG 임포트는 Airflow 3.0에서 도입된 안정적인 airflow.sdk 네임스페이스를 사용합니다.
Airflow Task SDK를 활용한 DAG 작성
Airflow 3.0에서 도입된 Task SDK는 DAG 정의를 Airflow 내부 구현으로부터 분리하는 독립 패키지입니다. 핵심 목표는 Airflow 업그레이드 시에도 코드 변경 없이 작동하는 이식 가능하고 안정적인 DAG를 작성하는 것입니다. DAG, dag, task, BaseOperator, Connection, Variable 등 모든 핵심 객체가 airflow.sdk 하위에 배치되어 있습니다.
레거시 임포트 경로(airflow.decorators.task, airflow.models.dag.DAG)는 3.2에서도 동작하지만 지원 중단으로 표시되어 있으며, 향후 릴리스에서 제거될 예정입니다.
# etl_daily_revenue.py
import pendulum
from airflow.sdk import dag, task
@dag(
schedule="@daily",
start_date=pendulum.datetime(2026, 1, 1, tz="UTC"),
catchup=False,
tags=["finance", "etl"],
)
def etl_daily_revenue():
@task()
def extract_transactions() -> list[dict]:
# Pulls raw transactions from the payments API
import requests
response = requests.get("https://api.internal/payments/daily")
return response.json()["transactions"]
@task(multiple_outputs=True)
def transform(transactions: list[dict]) -> dict:
# Aggregates by currency and computes totals
totals = {}
for tx in transactions:
currency = tx["currency"]
totals[currency] = totals.get(currency, 0) + tx["amount"]
return {"totals": totals, "count": len(transactions)}
@task()
def load(totals: dict, count: int):
# Inserts aggregated row into the warehouse
from warehouse import insert_revenue_summary
insert_revenue_summary(totals=totals, transaction_count=count)
raw = extract_transactions()
result = transform(raw)
load(result["totals"], result["count"])
etl_daily_revenue()@dag 데코레이터는 기존의 with DAG(...) 컨텍스트 매니저를 대체합니다. @task 데코레이터는 일반 Python 함수를 Airflow 태스크로 변환하며, XCom 직렬화는 자동으로 처리됩니다. transform(raw)를 호출하면 의존 관계가 선언되고, Airflow는 함수 호출 그래프를 기반으로 DAG 엣지를 구성합니다.
가변 워크로드를 위한 동적 태스크 매핑
처리할 항목 수가 실행마다 달라지는 경우, 정적 DAG로는 대응할 수 없습니다. Airflow 2.4에서 안정화되고 3.x에서 더욱 개선된 동적 태스크 매핑은 .expand()를 사용하여 런타임에 태스크를 확장함으로써 이 문제를 해결합니다.
# etl_multi_region.py
import pendulum
from airflow.sdk import dag, task
@dag(
schedule="@hourly",
start_date=pendulum.datetime(2026, 1, 1, tz="UTC"),
catchup=False,
tags=["regional", "etl"],
)
def etl_multi_region():
@task()
def get_regions() -> list[str]:
# Returns the active regions from config
return ["us-east-1", "eu-west-1", "ap-southeast-1"]
@task()
def process_region(region: str) -> dict:
# Processes events for a single region
from data_processor import aggregate_events
return aggregate_events(region)
@task()
def merge_results(results: list[dict]):
# Combines all regional outputs into a single report
from reporter import publish_global_report
publish_global_report(results)
regions = get_regions()
# .expand() creates one mapped task instance per region
processed = process_region.expand(region=regions)
merge_results(processed)
etl_multi_region()실행 시점에 Airflow는 리전별로 process_region의 병렬 인스턴스를 3개 생성합니다. 다음 날 네 번째 리전이 추가되더라도 DAG 코드를 수정할 필요가 없습니다. merge_results 태스크는 모든 매핑된 인스턴스가 완료될 때까지 대기합니다.
기본적으로 Airflow는 매핑당 1024개의 인스턴스로 제한합니다. 더 큰 데이터셋을 처리하는 경우 DAG 설정에서 max_map_length를 오버라이드하여 이 값을 변경할 수 있습니다.
에셋 파티션: Airflow 3.2의 데이터 인식 스케줄링
Airflow 3.2 이전에는 데이터 인식 스케줄링이 에셋 수준에서 작동했습니다. 프로듀서 DAG가 에셋을 업데이트하면, 어떤 데이터 슬라이스가 변경되었는지와 관계없이 모든 컨슈머 DAG가 트리거되었습니다. 에셋 파티션은 파티션 수준의 세분성을 지원하여 이 문제를 해결합니다.
세 개의 상위 DAG가 서로 다른 스포츠 리그의 시간별 선수 통계를 생성하는 시나리오를 생각해 봅시다. 하위 분석 DAG는 세 리그 모두 동일한 시간대의 데이터를 게시한 경우에만 트리거되어야 합니다.
# downstream_analytics.py
import pendulum
from airflow.sdk import dag, task
from airflow.timetables.assets import CronPartitionTimetable
# Define partitioned assets with hourly granularity
nba_stats = Asset("s3://datalake/nba/hourly/", partitions=CronPartitionTimetable("0 * * * *"))
epl_stats = Asset("s3://datalake/epl/hourly/", partitions=CronPartitionTimetable("0 * * * *"))
nfl_stats = Asset("s3://datalake/nfl/hourly/", partitions=CronPartitionTimetable("0 * * * *"))
@dag(
# Triggers only when all three assets have matching partition
schedule=(nba_stats & epl_stats & nfl_stats),
start_date=pendulum.datetime(2026, 1, 1, tz="UTC"),
catchup=False,
)
def unified_sports_analytics():
@task()
def aggregate_all_leagues(**context):
partition = context["partition"]
# Process only the specific hourly slice across all leagues
from analytics import build_cross_league_report
build_cross_league_report(partition=partition)
aggregate_all_leagues()
unified_sports_analytics()파티션 기반 스케줄링은 불필요한 파이프라인 실행을 제거합니다. 하위 DAG는 필요한 파티션이 모든 상위 소스에서 준비된 경우에만 실행되며, 부분적인 데이터로는 실행되지 않습니다.
Data Engineering 면접 준비가 되셨나요?
인터랙티브 시뮬레이터, flashcards, 기술 테스트로 연습하세요.
I/O 집약적 워크로드를 위한 네이티브 비동기 태스크
Airflow 3.2에서는 PythonOperator에 네이티브 비동기 지원이 추가되었습니다. 이전에는 대량의 API 호출이나 배치 파일 다운로드 같은 병렬 I/O 작업을 수행하려면 커스텀 deferrable 오퍼레이터를 작성해야 했습니다. 이제는 async 함수를 직접 전달할 수 있습니다.
# async_file_download.py
import asyncio
import aiohttp
from airflow.sdk import dag, task
import pendulum
@dag(
schedule="@daily",
start_date=pendulum.datetime(2026, 1, 1, tz="UTC"),
catchup=False,
tags=["async", "download"],
)
def async_file_download():
@task()
async def download_files():
urls = [f"https://data.source/files/{i}.csv" for i in range(500)]
async with aiohttp.ClientSession() as session:
# Downloads 500 files concurrently instead of sequentially
tasks = [fetch_and_save(session, url) for url in urls]
results = await asyncio.gather(*tasks)
return {"downloaded": len(results)}
download_files()
async def fetch_and_save(session: aiohttp.ClientSession, url: str):
async with session.get(url) as resp:
content = await resp.read()
filename = url.split("/")[-1]
with open(f"/data/downloads/{filename}", "wb") as f:
f.write(content)
return filename
async_file_download()비동기 방식은 단일 워커 슬롯에서 500개의 파일을 동시에 다운로드합니다. 동기 버전의 500회 순차 HTTP 요청과 비교하면, I/O 바운드 워크로드에서 수십 배의 속도 향상을 실현할 수 있으며 추가 워커를 프로비저닝할 필요가 없습니다.
Airflow 아키텍처: 구성 요소와 실행 흐름
Airflow의 아키텍처를 이해하는 것은 프로덕션 운영과 면접 대비 모두에 필수적입니다. 모든 DAG 실행에서 다섯 가지 구성 요소가 상호작용합니다.
- 스케줄러 — DAG 파일을 파싱하고, 의존 관계를 해석하며, 실행을 위해 태스크를 큐에 넣습니다. Airflow 3.x에서 스케줄러는 메타데이터 데이터베이스에 대해 무상태로 작동합니다.
- 익스큐터 — 태스크가 실행되는 위치를 결정합니다.
LocalExecutor는 단일 머신 환경을 처리합니다.CeleryExecutor는 워커 노드 간에 분산합니다.KubernetesExecutor는 태스크마다 파드를 생성합니다. - 워커 — 실제 태스크 코드를 실행합니다. Airflow 3.0의 Task Execution API를 통해 워커는 안정된 계약 기반으로 통신하며, 컨테이너, 엣지 환경, 외부 런타임에서의 실행이 가능합니다.
- 메타데이터 데이터베이스 — PostgreSQL(권장) 또는 MySQL을 사용합니다. DAG 정의, 태스크 상태, XCom 값, 연결 정보, 감사 로그를 저장합니다.
- 웹 서버 — DAG 실행 모니터링, 로그 확인, 수동 실행 트리거, 연결 관리를 위한 Airflow UI입니다.
프로덕션 환경에서는 SequentialExecutor 사용을 피해야 합니다. 한 번에 하나의 태스크만 실행할 수 있으며 개발 용도로만 존재합니다. Kubernetes 네이티브 환경에서는 KubernetesExecutor가 가장 강력한 격리를 제공합니다. 각 태스크가 독립적인 리소스와 의존성을 가진 개별 파드에서 실행되기 때문입니다.
2026년 Airflow vs. Prefect vs. Dagster 비교
Airflow에는 두 가지 주요 경쟁 도구가 있습니다. 올바른 선택은 팀 규모, 기존 인프라, 구축하는 파이프라인의 유형에 따라 달라집니다.
| 기능 | Airflow 3.2 | Prefect 3.x | Dagster 1.9 |
|---|---|---|---|
| DAG 정의 | Python 데코레이터(airflow.sdk) | Python 데코레이터(@flow, @task) | Python 데코레이터(@asset, @op) |
| 스케줄링 | Cron, 에셋 인식, 파티션 인식 | Cron, 이벤트 기반 | Cron, 센서 기반, 에셋 인식 |
| 실행 모델 | 중앙 스케줄러 + 분산 워커 | 하이브리드(서버 + 워크 풀) | 중앙 dagster-daemon |
| 동적 태스크 | .expand() 매핑된 태스크 | 네이티브 Python 루프 | 동적 파티션 |
| 비동기 지원 | 3.2에서 네이티브 지원 | 2.0부터 네이티브 지원 | 비동기 I/O 오퍼레이션 |
| 멀티팀 격리 | 내장(3.2 실험적 기능) | 워크스페이스 기반(Cloud) | 브랜치 배포 |
| 커뮤니티 규모 | 최대(GitHub 35,000+ 스타) | 성장 중(18,000+ 스타) | 성장 중(12,000+ 스타) |
| 적합한 환경 | 대규모 복잡한 멀티팀 파이프라인 | 빠른 개발 사이클을 원하는 소규모 팀 | 데이터 에셋 중심 조직 |
Airflow의 강점은 에코시스템에 있습니다. 주요 클라우드 서비스, 데이터베이스, API를 포괄하는 80개 이상의 프로바이더 패키지를 사용할 수 있습니다. Prefect는 보일러플레이트가 적은 우수한 개발자 경험에 강점이 있습니다. Dagster의 에셋 중심 모델은 태스크 시퀀스가 아닌 데이터 프로덕트 관점에서 사고하는 팀에 적합합니다.
데이터 엔지니어 면접에서 출제되는 Apache Airflow 질문
다음 질문들은 Airflow를 프로덕션 환경에서 운영하는 기업의 데이터 엔지니어 면접에서 실제로 출제되는 내용을 반영합니다.
DAG란 무엇이며, Airflow에서 어떻게 활용되는가
DAG(Directed Acyclic Graph, 방향 비순환 그래프)는 순환 참조가 없는 것이 보장된 태스크와 의존 관계의 집합으로 워크플로를 정의합니다. Airflow는 dags/ 폴더의 Python 파일을 파싱하여 의존 관계 그래프를 구축하고, 스케줄러가 실행 순서를 결정합니다. 각 DAG 실행은 논리적 날짜에 연결된 DagRun 객체를 생성합니다. "비순환" 제약 조건으로 인해 스케줄러는 항상 유효한 실행 순서를 찾을 수 있습니다.
XCom의 작동 방식과 사용을 피해야 하는 경우는 무엇인가
XCom(Cross-Communication)은 태스크 간에 소량의 데이터를 전달하는 메커니즘입니다. 태스크가 반환값을 XCom에 푸시하고, 하위 태스크가 이를 풀합니다. Task SDK에서는 데코레이팅된 태스크 간 함수 반환값 전달 시 자동으로 처리됩니다. XCom은 기본적으로 메타데이터 데이터베이스에 데이터를 저장하므로, 대용량 데이터(수 KB 이상)의 전송에는 적합하지 않습니다. 큰 데이터를 전송할 때는 외부 스토리지(S3, GCS)를 사용하고, XCom에는 참조 경로만 전달해야 합니다.
schedule, start_date, catchup의 차이점을 설명하라
schedule 파라미터(Airflow 3.x에서 schedule_interval에서 이름 변경)는 DAG의 실행 빈도를 정의합니다. cron 문자열, timedelta, 타임테이블 객체, 또는 에셋 트리거를 지정할 수 있습니다. start_date는 DAG 실행을 생성할 수 있는 가장 이른 논리적 날짜를 설정합니다. catchup=True(기본값)는 start_date부터 현재까지 누락된 모든 간격에 대해 DAG 실행을 생성합니다. catchup=False를 설정하면 과거 간격을 건너뛰고 현재 시점부터만 스케줄링합니다. 프로덕션의 일반적인 패턴은 운영 DAG에 catchup=False를, 이력 백필에 catchup=True를 설정하는 것입니다.
KubernetesExecutor와 CeleryExecutor의 차이점은 무엇인가
CeleryExecutor는 메시지 브로커(Redis 또는 RabbitMQ)를 통해 연결된 장시간 실행 워커 프로세스 풀을 유지합니다. 태스크가 큐에 들어가고 사용 가능한 워커에서 실행됩니다. KubernetesExecutor는 태스크마다 Docker 이미지와 리소스 요구사항을 기반으로 새로운 Kubernetes 파드를 생성합니다. Celery는 낮은 지연 시간(파드 시작 오버헤드 없음)을 제공하며 동질적 워크로드에 적합합니다. KubernetesExecutor는 더 강력한 격리와 태스크별 리소스 제어를 제공하며, 정적 워커 풀 관리가 불필요하여 이질적 워크로드에 이상적입니다.
태스크 실패 처리에는 어떤 전략이 있는가
Airflow는 여러 가지 실패 처리 메커니즘을 제공합니다. retries와 retry_delay로 지수 백오프를 적용한 자동 재시도를 설정할 수 있습니다. on_failure_callback은 태스크 실패 시 커스텀 로직(Slack 알림, PagerDuty 인시던트)을 트리거합니다. trigger_rule로 하위 태스크의 응답 방식을 제어합니다. all_success(기본값), one_success, all_failed, none_failed_min_one_success를 사용할 수 있습니다. 일시적인 인프라 장애에 대해서는 retry_exponential_backoff=True 파라미터로 재시도 간 대기 시간을 점진적으로 늘립니다. SLA(3.x에서는 Deadline Alerts)는 실행 시간을 모니터링하고, 태스크가 예상 실행 시간을 초과하면 콜백을 발동합니다.
Airflow 프로덕션 배포 모범 사례
Airflow를 대규모로 안정적으로 운영하려면 올바른 DAG 작성을 넘어서 여러 운영 패턴에 주의를 기울여야 합니다.
멱등성을 갖춘 태스크. 모든 태스크는 동일한 입력으로 여러 번 실행해도 같은 결과를 생성해야 합니다. 단순 INSERT 대신 INSERT ... ON CONFLICT나 MERGE를 사용합니다. 출력 데이터를 논리적 날짜 기준으로 파티셔닝합니다. 이를 통해 데이터 중복 없이 안전한 재시도와 백필이 가능해집니다.
작고 집중된 DAG. 수십 개의 태스크를 가진 모놀리식 DAG 구축을 지양해야 합니다. 복잡한 파이프라인은 에셋(이전 datasets)으로 연결된 여러 DAG로 분할합니다. 작은 DAG는 파싱이 빠르고 디버깅이 쉬우며, 부분적인 파이프라인 재시작이 가능합니다.
연결 관리. 모든 인증 정보를 Airflow의 연결 관리자 또는 외부 시크릿 백엔드(AWS Secrets Manager, HashiCorp Vault)에 저장합니다. DAG 파일에 인증 정보를 하드코딩해서는 안 됩니다. Airflow 3.x는 메타데이터 데이터베이스의 연결 필드를 저장 시 암호화합니다.
모니터링 및 알림. Airflow 메트릭을 Prometheus 또는 StatsD로 내보냅니다. scheduler_heartbeat, dag_processing.total_parse_time, executor.queued_tasks를 추적합니다. Airflow 3.2에서는 OpenTelemetry 트레이스가 추가되어 파이프라인의 엔드투엔드 관측 가능성이 확보되었습니다.
파이프라인 오케스트레이션을 포함한 데이터 엔지니어 면접을 준비하고 계신다면, SharpSkill의 Airflow 면접 대비 모듈을 활용할 수 있습니다. 프로덕션 환경의 실제 시나리오를 기반으로 한 연습 문제가 준비되어 있습니다.
연습을 시작하세요!
면접 시뮬레이터와 기술 테스트로 지식을 테스트하세요.
결론
- Airflow 3.2는 DAG 작성을 위한 안정 API로 Task SDK(
airflow.sdk)를 도입했습니다. 향후 호환성 문제를 방지하기 위해 임포트 마이그레이션을 권장합니다 - 에셋 파티션은 파티션 수준의 데이터 인식 스케줄링을 지원하여, 부분적 데이터 업데이트로 인한 불필요한 파이프라인 트리거를 제거합니다
- PythonOperator의 네이티브 비동기 지원으로 커스텀 deferrable 오퍼레이터 없이 I/O 집약적 워크로드를 처리할 수 있습니다
.expand()를 활용한 동적 태스크 매핑은 DAG 코드 변경 없이 가변적인 워크로드 크기에 런타임에서 대응합니다- 면접 준비 시 DAG 메커니즘, 익스큐터 트레이드오프(Celery vs. Kubernetes), XCom 제한사항, 멱등성 패턴을 반드시 숙지해야 합니다
- 프로덕션 배포에서는 작고 멱등한 DAG, 외부 시크릿 관리, 관측 가능성 플랫폼으로의 메트릭 내보내기가 효과적입니다
연습을 시작하세요!
면접 시뮬레이터와 기술 테스트로 지식을 테스트하세요.
공유
관련 기사

2026년 Delta Lake vs Apache Iceberg: 레이크하우스 아키텍처와 면접 대비 가이드
Delta Lake와 Apache Iceberg의 기술적 차이점을 상세히 분석합니다. 파티션 진화, ACID 트랜잭션, 쿼리 엔진 호환성 등 데이터 레이크하우스 면접에서 자주 출제되는 주제를 포괄적으로 다룹니다.

2026년 Snowflake: 아키텍처, SQL, 데이터 엔지니어 면접 질문
데이터 엔지니어를 위한 2026년 Snowflake 아키텍처 가이드입니다. 스토리지와 컴퓨트가 어떻게 분리되는지, virtual warehouse와 micro-partition이 어떻게 작동하는지, 그리고 실제 프로덕션 경험을 검증하는 면접 질문을 다룹니다.

dbt 2026 완벽 가이드: 데이터 변환, 테스트 전략, 면접 질문 총정리
dbt를 활용한 데이터 변환의 핵심 개념부터 실무까지, 레이어드 모델링, 인크리멘탈 전략, 테스트 방법론, 그리고 2026년 데이터 엔지니어링 면접에서 자주 출제되는 질문을 코드 예제와 함께 상세히 다룹니다.