왜 Pub/Sub와 List로는 부족한가

Redis의 Pub/Sub는 fire-and-forget 방식이다. 구독자가 순간적으로 끊겨 있으면 그 사이 발행된 메시지는 그냥 사라진다. LPUSH/BRPOP 기반의 List 큐는 메시지를 남기긴 하지만, 컨슈머가 메시지를 꺼낸 직후 처리 도중 죽으면 그 메시지는 유실된다. "꺼냈지만 아직 처리되지 않은" 상태를 추적할 방법이 없기 때문이다. 또한 여러 워커가 같은 큐를 나눠 처리하는 컨슈머 그룹 개념도 없다.

Redis Streams는 이 두 가지 문제를 정면으로 해결한다. 메시지는 로그처럼 append-only로 쌓이고, 컨슈머 그룹이 각 메시지의 소유권과 확인(ack) 상태를 서버 측에서 관리한다.

핵심 개념: 컨슈머 그룹과 PEL

스트림에 컨슈머 그룹을 만들면, 각 메시지는 그룹 내에서 정확히 한 컨슈머에게 배달된다. 컨슈머가 XREADGROUP으로 메시지를 읽는 순간 그 메시지는 PEL(Pending Entries List)에 등록된다. PEL은 "배달되었지만 아직 ack되지 않은" 메시지 목록이다. 컨슈머가 처리를 끝내면 XACK으로 PEL에서 제거한다. 만약 ack 없이 컨슈머가 죽으면 해당 메시지는 PEL에 남아 있으므로, 다른 컨슈머가 나중에 회수해 재처리할 수 있다. 이것이 at-least-once 보장의 근간이다.

그룹 생성과 발행

redis-cli XGROUP CREATE orders g1 $ MKSTREAM

redis-cli XADD orders '*' event_type paid order_id 1042 amount 39000

$는 "지금 이후에 도착하는 메시지부터"를 의미하고, MKSTREAM은 스트림이 없으면 만든다. XADD의 *는 서버가 <밀리초>-<시퀀스> 형식의 ID를 자동 부여하게 한다.

컨슈머 루프 구현

import redis

r = redis.Redis(decode_responses=True)
STREAM, GROUP, NAME = "orders", "g1", "worker-1"

while True:
    resp = r.xreadgroup(GROUP, NAME, {STREAM: ">"},
                        count=10, block=5000)
    if not resp:
        continue
    for _, messages in resp:
        for msg_id, fields in messages:
            try:
                handle(fields)          # 실제 비즈니스 처리
                r.xack(STREAM, GROUP, msg_id)
            except Exception:
                # ack 하지 않음 → PEL에 남아 재처리 대상
                log_failure(msg_id, fields)

>는 "아직 이 그룹의 누구에게도 배달되지 않은 새 메시지"를 뜻한다. block=5000으로 최대 5초간 대기하며 폴링 부하를 줄인다. 처리 성공 시에만 ack하는 것이 핵심이다.

죽은 컨슈머의 메시지 회수

워커가 죽으면 그 메시지는 PEL에 영원히 남는다. 이를 회수하려면 별도의 감시 루틴이 XAUTOCLAIM으로 오래된 pending 메시지를 자기 소유로 가져와 재처리한다.

# 60초 이상 처리되지 않은 메시지를 worker-1이 회수
redis-cli XAUTOCLAIM orders g1 worker-1 60000 0

min-idle-time(60000ms)은 "이만큼 방치된 것만 가져온다"는 뜻이라, 정상 처리 중인 메시지를 뺏지 않는다. 회수 후에도 반복 실패하는 메시지는 delivery count가 계속 오르므로, 임계값을 넘으면 dead-letter 스트림으로 옮겨 무한 재시도를 끊는다.

멱등성과 중복 처리

at-least-once는 곧 중복 배달 가능성을 의미한다. ack 직전에 죽으면 같은 메시지가 다시 처리된다. 따라서 컨슈머의 handle()은 반드시 멱등해야 한다. 예를 들어 메시지 ID나 비즈니스 키를 유니크 제약으로 걸어 두는 방식이다.

INSERT INTO processed_orders (order_id, amount)
VALUES (1042, 39000)
ON CONFLICT (order_id) DO NOTHING;

방식 비교

특성Pub/SubList 큐Streams + 그룹
오프라인 중 메시지 보존불가가능가능
처리 실패 시 재처리불가불가가능(PEL)
컨슈머 그룹 분산없음수동내장
전달 보장at-most-once사실상 없음at-least-once

운영 시 주의점

스트림은 append-only라 방치하면 메모리가 계속 증가한다. XADD에 MAXLEN ~ 100000 옵션을 붙이거나 주기적으로 XTRIM MINID로 오래된 항목을 잘라야 한다. 단, 아직 ack되지 않은 pending 메시지를 트림하지 않도록 보존 범위를 넉넉히 잡는다. 또한 컨슈머 이름은 워커마다 고정·유일하게 부여해야 XAUTOCLAIM 회수 로직이 정상 동작한다. 마지막으로 Redis 자체의 내구성은 RDB/AOF 설정에 달려 있으므로, 진짜 유실을 막으려면 AOF(everysec 이상)와 복제를 함께 고려해야 한다.