Celery Worker — 내부 구조와 동시성 모델

Celery Worker가 태스크를 어떻게 가져와 실행하는지, prefork·eventlet·gevent 동시성 모델의 차이, Prefetch 최적화, 태스크 생명주기와 재시도 전략을 설명한다.

Seobway · · 13분

Worker란 무엇인가

Celery Worker는 브로커에서 메시지를 가져와 태스크를 실행하는 프로세스다.[1]

Django + Celery 개요에서 설명한 전체 구조에서, Worker는 실제 일을 하는 주체다.
브로커(Redis/RabbitMQ)를 지속적으로 폴링(또는 PUSH 수신)해 새 태스크가 오면 꺼내 실행한다.

%% desc: Worker 내부 구조 — 메인 프로세스, Pool, 큐 소비자의 관계
flowchart TD
  subgraph WORKER["celery worker 프로세스 (메인)"]
    CTRL["컨트롤러\n(이벤트 루프)"]
    CONSUMER["큐 소비자\n(Broker 연결)"]

    subgraph POOL["실행 풀 (Pool)"]
      P1["자식 프로세스 1"]
      P2["자식 프로세스 2"]
      P3["자식 프로세스 3"]
      P4["자식 프로세스 4"]
    end
  end

  BROKER[메시지 브로커] --> CONSUMER
  CONSUMER -- "태스크 메시지" --> CTRL
  CTRL --> P1 & P2 & P3 & P4
  P1 & P2 & P3 & P4 --> RESULT[Result Backend]

동시성 모델

--concurrency 옵션으로 동시에 실행할 태스크 수를 지정한다.
기본값은 CPU 코어 수다.

Celery는 여러 동시성 모델을 지원한다.[2]

1. Prefork (기본값)

celery -A myproject worker --pool=prefork --concurrency=4
%% desc: Prefork 모델 — 각 자식 프로세스가 독립 메모리로 태스크를 병렬 처리
flowchart LR
  MAIN["메인 Worker"]
  P1["자식 1\nPID 1001\n독립 메모리"]
  P2["자식 2\nPID 1002\n독립 메모리"]
  P3["자식 3\nPID 1003\n독립 메모리"]

  MAIN --> P1
  MAIN --> P2
  MAIN --> P3

  P1 --> T1["태스크 A\nCPU 연산"]
  P2 --> T2["태스크 B\nDB 처리"]
  P3 --> T3["태스크 C\nPDF 생성"]

2. Eventlet / Gevent

pip install eventlet
celery -A myproject worker --pool=eventlet --concurrency=100

3. Threads

celery -A myproject worker --pool=threads --concurrency=8

동시성 모델 선택 기준

모델 적합한 태스크 동시성 수준 메모리
Prefork CPU 집약적 (이미지 처리, ML) CPU 코어 수 높음
Eventlet I/O 집약적 (API 호출, 이메일) 수백 낮음
Gevent I/O 집약적 수백 낮음
Threads I/O 집약적 수십 중간

Prefetch — 미리 가져오기

Worker는 브로커에서 태스크를 하나씩 가져오지 않고, 미리 여러 개를 가져온다.
이를 Prefetch라 한다.

# settings.py
CELERY_WORKER_PREFETCH_MULTIPLIER = 1   # 기본값: 4

기본 설정(prefetch_multiplier=4, concurrency=4)이면:

문제: 태스크 실행 시간이 불균일할 때, 빠른 Worker가 놀고 느린 Worker가 16개를 독점할 수 있다.

권장 설정 — 긴 태스크가 있는 경우:

CELERY_WORKER_PREFETCH_MULTIPLIER = 1
CELERY_TASK_ACKS_LATE = True         # 실행 완료 후 ACK
%% desc: Prefetch 1 vs 4 비교 — Prefetch 낮을수록 태스크가 균등하게 분배됨
flowchart LR
  subgraph A["Prefetch=4 (기본)"]
    B1["큐: 16개 태스크"] --> W1A["Worker A\n12개 점유"]
    B1 --> W2A["Worker B\n4개 점유"]
    B1 --> W3A["Worker C\n0개 → 유휴"]
  end

  subgraph B["Prefetch=1 (권장)"]
    B2["큐: 16개 태스크"] --> W1B["Worker A\n4개 균등"]
    B2 --> W2B["Worker B\n4개 균등"]
    B2 --> W3B["Worker C\n4개 균등"]
  end

태스크 생명주기

%% desc: Celery 태스크의 전체 생명주기 — 발행부터 결과 저장까지
sequenceDiagram
  participant PROD as Producer (Django)
  participant BROKER as 브로커 (Redis)
  participant WORKER as Worker
  participant BACKEND as Result Backend

  PROD->>BROKER: task.delay(args) → 메시지 발행
  Note over BROKER: 큐에 저장 (PENDING 상태)

  WORKER->>BROKER: 메시지 소비 (ACK 전)
  Note over WORKER: STARTED 상태
  WORKER->>WORKER: 태스크 함수 실행

  alt 성공
    WORKER->>BACKEND: 결과값 저장
    WORKER->>BROKER: ACK 전송 (메시지 제거)
    Note over BACKEND: SUCCESS 상태
  else 실패 + 재시도
    WORKER->>BROKER: 재시도 메시지 발행
    Note over WORKER: RETRY 상태
  else 최종 실패
    WORKER->>BACKEND: 예외 정보 저장
    WORKER->>BROKER: ACK 전송
    Note over BACKEND: FAILURE 상태
  end

태스크 재시도 전략

자동 재시도 (autoretry_for)

@shared_task(
    bind=True,
    autoretry_for=(requests.exceptions.RequestException,),
    retry_kwargs={"max_retries": 5},
    retry_backoff=True,          # 지수 백오프
    retry_backoff_max=300,       # 최대 5분 대기
    retry_jitter=True,           # 랜덤 지터 추가 (동시 재시도 분산)
)
def fetch_external_data(self, url: str):
    response = requests.get(url, timeout=30)
    response.raise_for_status()
    return response.json()

재시도 대기 시간 (지수 백오프 + jitter):

수동 재시도 (self.retry)

@shared_task(bind=True)
def send_sms(self, phone: str, message: str):
    try:
        sms_client.send(phone, message)
    except SmsQuotaExceeded:
        # 1시간 후 재시도
        raise self.retry(exc=SmsQuotaExceeded(), countdown=3600)
    except SmsInvalidNumber:
        # 재시도 불필요한 오류 → 그냥 실패
        raise

태스크 라우팅 — 큐 분리

태스크 종류별로 다른 큐, 다른 Worker를 사용할 수 있다.

# settings.py
CELERY_TASK_ROUTES = {
    "myapp.tasks.send_email": {"queue": "email"},
    "myapp.tasks.generate_report": {"queue": "reports"},
    "myapp.tasks.*": {"queue": "default"},
}
# 이메일 전용 Worker (Eventlet, 고동시성)
celery -A myproject worker -Q email --pool=eventlet --concurrency=100

# 리포트 전용 Worker (Prefork, CPU 집약)
celery -A myproject worker -Q reports --pool=prefork --concurrency=2

# 기본 Worker
celery -A myproject worker -Q default --concurrency=4

Worker 모니터링 — Flower

pip install flower
celery -A myproject flower --port=5555

Flower는 Web UI로 Worker 상태, 태스크 히스토리, 실시간 처리율을 볼 수 있다.[3]

참고

  1. Celery Workers Guide, Celery Docs
  2. Celery Concurrency, Celery Docs
  3. Flower — Celery monitoring tool, Flower Docs

관련 글