2026년 8월 15일
[1편] Terraform+Docker로 로컬 ELT 파이프라인 만들기 — NYC Taxi 데이터로 Airflow 완주하기
프로젝트 목표
클라우드 리소스 없이 로컬 환경에서 Terraform, Docker, Airflow, Celery, Redis, MinIO, DuckDB, dbt 여덟 개 기술로 ELT 파이프라인을 한 바퀴 완주하는 것이 이 프로젝트의 목표다. 실무에서 각 기술이 파이프라인 안에서 어떤 역할을 맡는지, 그리고 그 기술들이 서로 어떻게 맞물리는지를 직접 손으로 만들어보며 확인하는 데 목적이 있다.
데이터셋은 TPC-H 대신 NYC Yellow Taxi의 월별 parquet 파일을 채택했다. TPC-H는 한 번에 통째로 생성되는 벤치마크용 데이터라서 "새 파일이 도착하면 그것만 적재한다"는 ELT의 증분 처리 패턴을 자연스럽게 재현하기 어렵다. 반면 월별로 나뉜 parquet 파일은 파일 단위·월 단위 증분 적재를 눈으로 확인하기에 적합하다.
전체 인프라를 처음부터 한꺼번에 띄우면 문제가 생겼을 때 원인 후보가 지나치게 많아진다. 따라서 이 프로젝트는 두 단계로 나뉜다. LocalExecutor로 먼저 Airflow의 기본 개념을 잡는 단계(Phase A, 이번 편)와, 이를 CeleryExecutor로 전환하는 단계(Phase B, 다음 편)다.
등장하는 기술의 역할
- Airflow: Extract/Load/Transform 태스크의 순서·의존성·스케줄을 관리하는 오케스트레이터다. "무엇을, 언제, 어떤 순서로 실행할지"만 결정하고, 실제 실행은 별도 주체(Executor)에게 맡긴다.
- MinIO: S3 호환 오브젝트 스토리지다. 이 프로젝트에서는 원본 parquet 파일이 가공 전 그대로 놓이는 data lake, 즉 Extract와 Load 사이의 중간 저장소 역할만 한다.
- DuckDB: 임베디드 OLAP 엔진이다. 별도 서버 없이 로컬 파일 하나로 동작하는 data warehouse로 사용하며, raw 테이블과 dbt가 만드는 mart 테이블이 이 안에 함께 존재한다.
- dbt: DuckDB 안에서 raw 테이블을 staging을 거쳐 mart로 변환하는 SQL 기반 변환 도구다. SQL을 소프트웨어처럼 버전관리하고 테스트한다는 것이 핵심 아이디어다.
- Celery/Redis: 이번 편(Phase A)에서는 등장하지 않는다. Airflow가 태스크를 누구에게 실행시킬지 정하는 계층(Executor)을 스케줄러 자신에서 Celery worker로 바꾸는 것이 다음 편의 주제다.
아키텍처
flowchart LR
subgraph E["Extract"]
E1["로컬 월별 parquet"] --> E2["MinIO 버킷 업로드"]
end
subgraph L["Load"]
L1["MinIO → 로컬 스테이징 다운로드<br/>→ DuckDB 적재 (증분)"]
end
subgraph T["Transform"]
T1["dbt run: staging → mart"]
end
E -->|"raw/taxi/2023-01.parquet"| L
L -->|"raw_trips 테이블"| TAirflow DAG 하나(nyc_taxi_elt)가 Extract → Load → Transform 세 태스크의 순서와 의존성을 관리한다. 이번 편에서는 이 태스크들을 스케줄러 자신이 직접 실행하는 LocalExecutor를 사용한다.
설계 결정 두 가지
1. DuckDB가 MinIO를 직접 쿼리하지 않는다. httpfs/S3 확장을 쓰면 DuckDB가 오브젝트 스토리지를 직접 읽게 만들 수도 있지만, 이 프로젝트에서는 Load 태스크(PythonOperator)가 minio SDK로 파일을 로컬로 내려받은 뒤 DuckDB가 그 로컬 경로만 보도록 분리했다. Extract(로컬→MinIO)와 Load(MinIO→로컬→DuckDB)를 물리적으로 나누면, MinIO는 중간 저장소 역할만 하면 되고 DuckDB는 S3 프로토콜을 몰라도 된다.
# Extract: 로컬 원본 -> MinIO 업로드
client.fput_object(bucket, object_name, local_path)
# Load: MinIO -> 로컬 스테이징 다운로드
client.fget_object("raw", object_name, stage_path)2. 파일 존재 확인은 버킷 리스팅 대신 메타 테이블로 한다. 매번 MinIO 버킷 목록과 DuckDB 내용을 비교하는 대신, DuckDB 안에 적재 이력을 남기는 테이블을 둔다. 실무의 워터마크·적재 매니페스트 패턴과 같은 개념이다.
CREATE TABLE IF NOT EXISTS loaded_files (
file_name VARCHAR PRIMARY KEY,
loaded_at TIMESTAMP
);already_loaded = con.execute(
"SELECT 1 FROM loaded_files WHERE file_name = ?", [file_name]
).fetchone()
if already_loaded:
print(f"{file_name} already loaded, skip.")
return이 두 결정 덕분에 같은 월을 여러 번 실행해도 중복 적재가 일어나지 않는 멱등성을 확인할 수 있었다.
DAG 구조
with DAG(
dag_id="nyc_taxi_elt",
schedule=None,
catchup=False,
params={"target_month": Param("2023-01", type="string")},
) as dag:
extract = PythonOperator(task_id="extract_to_minio", python_callable=extract_to_minio)
load = PythonOperator(task_id="load_to_duckdb", python_callable=load_to_duckdb)
transform = BashOperator(
task_id="dbt_run",
bash_command=f"cd {DBT_PROJECT_DIR} && DBT_PROFILES_DIR={DBT_PROJECT_DIR} dbt run",
)
extract >> load >> transform대상 월은 하드코딩하지 않고 Airflow의 dag_run.conf 파라미터(target_month)로 받는다. UI에서 매번 다른 달을 입력해 트리거할 수 있다. 세 태스크는 >> 연산자로 의존관계만 선언하며, 실행 순서·재시도·로그 수집은 Airflow가 담당한다.
Transform — dbt로 raw를 staging/mart로
마지막 태스크(dbt_run)가 실제로 하는 일은 raw_trips를 그대로 둔 채 그 위에 정리된 모델을 쌓는 것이다.
-- models/staging/stg_trips.sql
select
tpep_pickup_datetime as pickup_at,
tpep_dropoff_datetime as dropoff_at,
trip_distance,
fare_amount,
passenger_count
from {{ source('raw', 'raw_trips') }}-- models/marts/mart_daily_trips.sql
select
date_trunc('day', pickup_at) as trip_date,
count(*) as trip_count,
avg(fare_amount) as avg_fare,
avg(trip_distance) as avg_distance
from {{ ref('stg_trips') }}
group by 1ref()와 source()로 모델 간 의존관계를 선언하면 dbt가 실행 순서를 알아서 정하고, dbt docs generate로 이 의존관계를 lineage 그래프로 확인할 수 있다. raw 테이블을 직접 건드리지 않고 staging에서 한 번 정리한 뒤 mart를 쌓는 이유는, 원본을 훼손 없이 보존하면서도 변환 로직을 반복 가능하게 만들기 위함이다.
트러블슈팅 — 로그가 안 보인다는 증상, 원인은 세 가지였다
Airflow는 태스크 로그를 스케줄러·워커가 로컬에 쓰고, 웹서버가 이를 HTTP(기본 8793 포트)로 해당 컨테이너에 요청해서 UI에 보여주는 구조다. 공식 docker-compose.yaml은 이 구조가 정상 동작하도록 볼륨 공유, SECRET_KEY 등을 이미 세팅해둔 상태인데, Terraform으로 직접 구성하면서 이 전제들을 하나씩 놓쳤다.
케이스 1 — Invalid URL 'http://:8793/...': No host supplied
scheduler와 webserver 컨테이너 간에 logs 디렉터리를 볼륨으로 공유하지 않은 것이 원인이었다. 양쪽 컨테이너에 동일한 host 경로를 /opt/airflow/logs로 마운트해 해결했다.
케이스 2 (오해였던 것) — 특정 태스크만 로그·hostname이 안 남음
확인용 DAG를 EmptyOperator로 만들었는데, Airflow는 EmptyOperator를 스케줄러가 실제 실행 없이 바로 success 처리하는 최적화를 갖고 있다. 버그가 아니라 애초에 로그가 안 남는 것이 정상 동작이다. 파이프라인 동작 확인용 더미 태스크가 필요하다면 PythonOperator처럼 실제로 실행되는 Operator를 써야 한다.
케이스 3 — 403 FORBIDDEN
AIRFLOW__WEBSERVER__SECRET_KEY를 고정하지 않아서, terraform apply로 컨테이너를 재생성할 때마다 webserver와 scheduler가 각자 다른 랜덤 secret_key를 생성한 것이 원인이었다. 로그 서버(8793) 인증(JWT)이 서로 맞지 않아 403이 발생했다. 케이스 1과는 완전히 다른 원인인데 증상(로그가 안 보임)만 보면 구분이 되지 않는다. local.airflow_env에 AIRFLOW__WEBSERVER__SECRET_KEY를 고정 문자열로 추가해 해결했다.
부가 — S3Error: NoSuchKey
트리거 방식(CLI vs DAG 파라미터)을 오가며 테스트하다가, 실제로는 Extract 태스크(파일 업로드)를 거치지 않고 Load 태스크부터 실행해서 발생한 에러였다. 인프라 문제가 아니라 DAG 의존 순서를 지키지 않은 사용자 실수였다.
비교 — 증상별 원인 정리
| 증상 | 원인 | 성격 |
|---|---|---|
| Invalid URL / host 없음 | logs 볼륨 미공유 | 설정 실수 |
| 특정 태스크만 로그 없음 | EmptyOperator 최적화 | 정상 동작 (오해) |
| 403 Forbidden | SECRET_KEY 미고정 → 재생성마다 값 바뀜 | 설정 실수 |
| S3Error NoSuchKey | Extract 생략하고 Load부터 실행 | 사용자 실수 (DAG 의존 순서) |
정리하면
- Airflow를 Terraform 등으로 직접 구성할 때는 공식 compose가 암묵적으로 해주는 것들(logs/dags 볼륨 공유, SECRET_KEY/FERNET_KEY 고정)을 먼저 체크리스트로 뽑고 시작해야 한다.
- SECRET_KEY류는 컨테이너 재생성마다 값이 바뀌면 증상이 매번 다르게(빈 URL, 403, connection refused) 나타나서 원인 특정에 시간이 오래 걸린다. 처음부터 고정값으로 박아두는 것이 낫다.
- 파이프라인 동작 확인용 더미 태스크는 EmptyOperator 대신 실제로 실행되는 Operator를 써야 한다.
- Extract와 Load를 물리적으로 분리하고 메타 테이블로 적재 이력을 관리하면, 별도 장치 없이도 재실행 시 멱등성이 자연스럽게 보장된다.
- raw와 mart를 dbt로 분리해두면, 변환 로직이 바뀌어도 원본을 다시 수집하지 않고 재실행할 수 있다.
다음 편 예고
이번 편은 Airflow가 태스크를 스케줄러 자신이 직접 실행하는 LocalExecutor로 완주했다. 다음 편에서는 Redis, Celery Worker, Triggerer를 추가하고 AIRFLOW__CORE__EXECUTOR만 LocalExecutor에서 CeleryExecutor로 바꾼다. DAG 코드는 한 줄도 바꾸지 않지만 실행 주체는 Celery worker로 넘어간다. "무엇을 실행할지 결정하는 계층"과 "실제로 실행하는 계층"이 분리되어 있어서 후자만 교체할 수 있다는 Executor 추상화를 직접 확인하는 것이 다음 편의 주제다.
참고
- NYC TLC Trip Record Data
- Airflow 공식
docker-compose.yaml(Running Airflow in Docker)