seongju
PORTFOLIO
He who has a why to live can bear almost any how. — NietzscheHe who has a why to live can bear almost any how. — NietzscheHe who has a why to live can bear almost any how. — NietzscheHe who has a why to live can bear almost any how. — Nietzsche

2026년 9월 1일

[2편] LocalExecutor에서 CeleryExecutor로 — DAG 코드는 그대로, 실행 계층만 교체

1편에서 이어서

1편에서는 Terraform + Docker로 Airflow(LocalExecutor) + MinIO + DuckDB + dbt를 로컬에 띄우고, NYC Yellow Taxi 월별 parquet을 nyc_taxi_elt DAG 하나로 Extract → Load → Transform 완주했다. 이때 태스크를 실제로 실행한 건 스케줄러 프로세스 자신이었다.

이번 편의 목표는 딱 하나다. DAG 코드를 한 줄도 바꾸지 않고, 태스크를 실행하는 주체를 스케줄러에서 별도의 Celery worker로 옮긴다. 바꾸는 것은 환경변수 세 개와 컨테이너 세 개뿐이다.

Executor라는 레이어

Airflow에서 스케줄러가 하는 일과 executor가 하는 일은 분리되어 있다.

  • 스케줄러: "무엇을, 언제, 어떤 순서로 실행할지"를 결정한다. DAG를 파싱하고, 의존관계가 풀린 태스크를 찾아 "실행 대기(queued)" 상태로 만든다.
  • Executor: queued된 태스크를 받아 "실제로 어디서 어떻게 실행할지"를 담당한다.

LocalExecutor는 스케줄러 프로세스가 태스크를 자식 프로세스로 직접 띄운다. 스케줄러가 죽으면 실행 중이던 태스크도 같이 죽고, 실행 용량은 그 한 대의 스케줄러 리소스에 묶인다.

CeleryExecutor는 태스크를 메시지 브로커에 밀어넣기만 하고, 브로커에서 태스크를 꺼내 실행하는 것은 완전히 별개의 프로세스(Celery worker)다. worker는 여러 대로 늘릴 수 있고, 스케줄러와 생명주기가 분리된다.

추가되는 컴포넌트

컴포넌트 역할
Redis 메시지 브로커. 스케줄러가 "이 태스크 실행해" 메시지를 넣는 큐. worker는 여기서 메시지를 꺼낸다.
Celery worker 브로커에서 태스크 메시지를 꺼내 실제로 실행하는 프로세스.
result backend 태스크가 어떻게 끝났는지(성공/실패)를 저장하는 곳. worker가 실행을 끝내면 여기에 결과를 쓰고, CeleryExecutor는 이걸 폴링해 태스크 상태를 추적한다.
Triggerer 이 프로젝트에선 안 쓰지만 Airflow 공식 docker-compose.yaml 구성의 표준 멤버라 같이 올린다. (특정 유형의 대기성 태스크를 스케줄러 밖에서 처리하는 별도 프로세스다.)

설계 결정 — 브로커와 result backend를 분리한다

AIRFLOW__CELERY__BROKER_URL은 Redis로, AIRFLOW__CELERY__RESULT_BACKENDdb+postgresql://...로 따로 잡았다. 1편 기술 목록엔 안 넣었지만, Airflow는 DAG·태스크 상태를 저장할 메타데이터 DB가 필요하고 이 프로젝트는 처음부터 Postgres 컨테이너를 그 용도로 띄워두고 있었다.

둘 다 Redis로 몰 수도 있다. 하지만 result backend는 나중에 조회되고 어느 정도 영속성이 필요한 데이터다. 이미 Postgres가 떠 있으니 거기에 얹는 편이 TTL·유실 관리 부담이 적다. Redis는 순수하게 "태스크를 흘려보내는 큐"로만 쓴다.

아키텍처

flowchart LR
    SCH["Scheduler<br/>(무엇을/언제/순서)"] -->|"태스크 메시지 push"| REDIS[("Redis<br/>broker")]
    REDIS -->|"태스크 메시지 pull"| WORKER["Celery Worker<br/>(실제 실행)"]
    WORKER -->|"Extract/Load"| MINIO[("MinIO")]
    WORKER -->|"Load/Transform"| DUCK[("DuckDB<br/>warehouse")]
    WORKER -->|"실행 결과"| PG[("Postgres<br/>result backend")]
    SCH --- PG
    WEB["Webserver"] -->|"공유 볼륨에서<br/>로그 직접 read"| LOGS[/"/opt/airflow/logs"/]
    WORKER --> LOGS

1편에서 로그 디렉터리(/opt/airflow/logs)를 스케줄러와 웹서버에 같은 경로로 물렸던 것(1편 트러블슈팅 케이스 1)이 여기서 다시 효과를 본다. worker도 같은 볼륨에 로그를 쓰므로, 웹서버는 1편에서 다룬 8793 로그 서버를 거치지 않고 파일을 직접 읽는다.

바뀌는 것 / 바뀌지 않는 것

바뀌는 것 — 환경변수 3개:

- AIRFLOW__CORE__EXECUTOR=LocalExecutor
+ AIRFLOW__CORE__EXECUTOR=CeleryExecutor
+ AIRFLOW__CELERY__BROKER_URL=redis://redis:6379/0
+ AIRFLOW__CELERY__RESULT_BACKEND=db+postgresql://airflow:***@postgres:5432/airflow

그리고 컨테이너 3개(redis, worker, triggerer) 추가.

바뀌지 않는 것: dags/nyc_taxi_elt.py는 한 줄도 안 바뀐다. dbt 프로젝트도, Airflow 이미지도(apache-airflow-providers-celery는 공식 이미지에 이미 포함), Extract → Load → Transform 로직과 >> 의존관계 선언도 그대로다.

Terraform — 스위치 하나

Phase A/B를 use_celery 변수 하나로 토글한다.

variable "use_celery" {
  type    = bool
  default = false
}

locals {
  executor = var.use_celery ? "CeleryExecutor" : "LocalExecutor"

  airflow_env = concat(
    local.base_airflow_env,          # EXECUTOR 값은 local.executor 사용
    var.use_celery ? local.celery_airflow_env : [],
  )
}

# redis / worker / triggerer 는 count 로 껐다 켠다
resource "docker_container" "airflow_worker" {
  count   = var.use_celery ? 1 : 0
  command = ["${local.wait_for_db} && exec airflow celery worker"]
  # ...
}
terraform apply -var use_celery=true

plan 결과는 3개 add(redis/worker/triggerer), 3개 replace(scheduler/webserver/init — 환경변수만 바뀌어서). DAG 파일도 dbt도 이미지도 건드리지 않는다.

검증 — 실행 host가 바뀌었는가

같은 DAG를 트리거하고, task_instance 테이블의 hostname 컬럼을 직접 봤다.

select dr.run_id, ti.task_id, ti.state, ti.hostname
from dag_run dr join task_instance ti using (dag_id, run_id)
where dr.dag_id = 'nyc_taxi_elt'
order by dr.start_date desc;
run task hostname 정체
LocalExecutor 시절 extract / load / dbt_run 491044066bb5 scheduler 컨테이너
CeleryExecutor 전환 후 extract / load / dbt_run f8914b104683 worker 컨테이너

DAG 코드는 동일한데 태스크가 실행된 컨테이너가 스케줄러에서 worker로 넘어갔다. "무엇을 실행할지 결정하는 계층"과 "실제로 실행하는 계층"이 분리되어 있어서 후자만 교체할 수 있다는 것이 눈으로 확인된다.

한계 — 공유 warehouse의 동시성

executor를 바꾸면 실행이 분산된다. 그런데 그 아래에서 모두가 쓰는 DuckDB 파일은 그 분산을 견디지 못한다.

증상. Load 태스크에서 이 에러가 난다.

IO Error: Could not set lock on file "/opt/airflow/duckdb/warehouse.duckdb":
Conflicting lock is held in /usr/local/bin/python3.12 (PID 359)

원인. DuckDB는 파일 단위 single-writer다. 두 프로세스가 같은 .duckdb 파일을 쓰기로 열면 뒤에 온 쪽이 이 에러로 튕긴다.

LocalExecutor에서는 이게 "여러 DAG run이 병렬로 떠서 각자의 Load가 동시에 warehouse를 열었다"일 때만 났다. 이 프로젝트의 DAG에는 max_active_runs=1을 걸어 DAG run을 한 번에 하나로 직렬화해뒀는데(1편 본문에선 따로 짚지 않았다), 그래서 평소엔 안 부딪힌다.

CeleryExecutor에서는 이 제약이 구조적으로 더 크게 다가온다. worker는 물리적으로 분리된 프로세스이고(스케일하면 여러 컨테이너), worker 하나 안에서도 Celery의 기본 동시성이 16이라 16개 프로세스가 브로커에서 태스크를 집어간다. warehouse에 동시 쓰기가 필요한 순간 "어느 worker의 어느 프로세스가 파일 락을 쥐고 있나"가 바로 문제가 된다.

대응. 두 가지를 같이 쓴다.

  1. max_active_runs=1로 warehouse 쓰기를 직렬화한다. 그리고 Load 태스크에 retries를 붙여, 그래도 겹쳐서 락에 튕기면 재시도로 흡수한다.
  2. Load를 idempotent하게 유지해서 재시도가 안전하게 만든다. 이미 1편에서 loaded_files 메타 테이블로 "같은 파일은 두 번 안 넣는다"를 구현해뒀으므로, 재시도가 중복 적재를 만들지 않는다.

로컬 학습용이라 여기까지다. 실무 규모라면 답은 하나다 — 애초에 동시 쓰기가 되는 warehouse(Postgres, Redshift, BigQuery 등)로 바꾼다. 즉 실행 계층을 분산시키는 것과, 그 아래 저장소가 그 분산을 견디는 것은 별개의 설계 문제다.

곁가지 — 소스 스키마는 티가 안 난다.

파이프라인을 여러 달 돌리다 알게 된 것 하나. NYC TLC의 2023년 Yellow Taxi parquet은 2023-01에는 airport_fee(소문자), 2023-02부터는 Airport_fee(대문자 A)로 컬럼명이 바뀌어 있다. 이 프로젝트의 Load(INSERT INTO raw_trips SELECT * FROM read_parquet(...))는 컬럼을 이름이 아니라 위치로 매칭하기 때문에 그냥 적재됐지만, 이름 기준이었다면 2023-02에서 깨졌을 것이다. 소스가 스키마를 바꿔도 파이프라인이 버티느냐는 이런 식으로 우연히 결정된다 — 이건 3편에서 dbt 테스트로 다시 다룬다.

정리하면

  • CeleryExecutor 전환은 AIRFLOW__CORE__EXECUTOR + broker/result-backend 환경변수 3개와 컨테이너 3개 추가가 전부다. DAG·dbt·이미지는 그대로다.
  • broker(Redis)와 result backend(Postgres)는 분리하는 편이 관리가 쉽다. Redis는 큐로만.
  • task_instance.hostname으로 실행 host가 scheduler → worker로 넘어간 것을 확인할 수 있다.
  • DAG는 "무엇을, 어떤 순서로"만 선언한다. >>는 의존관계이지 실행 방법이 아니다. "누가, 어디서, 몇 개까지 동시에"는 executor 레이어가 캡슐화하므로, DAG를 그대로 두고 executor만 교체할 수 있다. 학습은 LocalExecutor로, 확장은 CeleryExecutor로 단계를 나눌 수 있는 이유다.
  • 실행 계층을 분산시켜도 그 아래 공유 자원(DuckDB 파일)의 동시성은 max_active_runs나 warehouse 교체로 따로 풀어야 한다.

1편에서 예고한 대로 Celery와 Redis가 붙으면서, "클라우드 없이 여덟 개 기술로 ELT를 한 바퀴" 라는 1편의 목표는 여기서 채워졌다. LocalExecutor로 개념을 잡고 CeleryExecutor로 실행 계층을 갈아끼우는 것까지가 여기까지의 이야기다.

그런데 파이프라인이 도는 것과 나온 데이터가 맞는 것은 다른 문제다. mart를 열어보면 2001년, 2008년 같은 말도 안 되는 날짜의 행이 섞여 있고, 요금이 음수인 행도 있다. 3편에서는 dbt의 test와 dbt build로 "이 데이터를 믿어도 되는가"를 파이프라인 안에 넣는다.

참고