Pub/Sub로 알림을 구현했는데 서버가 잠깐 내려가면 그동안의 메시지가 사라집니다. Redis Streams는 이 문제를 해결합니다. 메시지를 저장하고, 누가 처리했는지 추적하며, 실패 시 재처리도 가능합니다.
Stream은 로그처럼 동작하는 자료구조입니다.
소비자 그룹(Consumer Group): Stream에서 여러 소비자가 메시지를 나눠 처리하는 단위. 같은 그룹 내에서는 메시지가 한 소비자에게만 전달됩니다.
# 메시지 추가 (* = 자동 ID 생성)
XADD orders * action "placed" userId 1000 amount 29000
# "1718700000000-0" 형태의 ID 반환
# ID는 타임스탬프-순번 형식
# 1718700000000 = Unix 밀리초 타임스탬프
# 0 = 같은 밀리초 내 순번
# 최대 크기 제한 (오래된 메시지 자동 삭제)
XADD orders MAXLEN 10000 * action "placed" userId 1001 amount 15000
# 처음부터 읽기 (0 = 시작)
XREAD COUNT 10 STREAMS orders 0
# 특정 ID 이후 읽기
XREAD COUNT 10 STREAMS orders 1718700000000-0
# 블로킹으로 대기 (0 = 무한 대기)
XREAD COUNT 1 BLOCK 0 STREAMS orders $
# $ = 이 명령 실행 이후 새로 추가되는 것만
# 전체 조회
XRANGE orders - +
# 특정 ID 범위
XRANGE orders 1718700000000-0 1718786400000-0
# 최근 10개
XREVRANGE orders + - COUNT 10
# 길이
XLEN orders
XRANGE: - 는 가장 작은 ID(처음), + 는 가장 큰 ID(끝)를 의미합니다. 시작과 끝 범위를 지정하면 그 사이 메시지를 조회합니다.
여러 Worker가 메시지를 나눠 처리하는 핵심 기능입니다.
# 그룹 생성 ($ = 이후 메시지부터, 0 = 처음부터)
XGROUP CREATE orders email-workers $
XGROUP CREATE orders analytics-workers 0
# 메시지 읽기 (> = 아직 전달 안 된 메시지)
XREADGROUP GROUP email-workers worker-1 COUNT 5 STREAMS orders >
# 처리 완료 확인 (ACK)
XACK orders email-workers 1718700000000-0
Stream: orders
───────────────────────────────────
1718700000000-0 (주문 A)
1718700001000-0 (주문 B)
1718700002000-0 (주문 C)
email-workers 그룹:
worker-1 → 주문 A 처리 중
worker-2 → 주문 B 처리 중
analytics-workers 그룹:
worker-3 → 주문 A, B, C 읽음
같은 Stream을 두 그룹이 독립적으로 소비합니다.
ACK(Acknowledgment): 소비자가 메시지를 처리했음을 확인하는 응답. ACK를 보내지 않으면 Redis는 해당 메시지를 미처리로 추적합니다.
ACK를 받지 못한 메시지는 PEL(Pending Entry List)에 남습니다.
# 미처리 메시지 목록 조회
XPENDING orders email-workers - + 10
# 오래된 미처리 메시지 다른 Worker에게 재할당
XCLAIM orders email-workers worker-2 60000 1718700000000-0
# 60000ms(1분) 이상 미처리인 메시지를 worker-2에게 재할당
PEL(Pending Entry List): 소비자 그룹에서 전달되었지만 아직 ACK를 받지 못한 메시지 목록. 실패한 Worker의 메시지를 다른 Worker가 인계받을 수 있습니다.
BullMQ는 내부적으로 Redis Streams를 사용합니다. BullMQ가 제공하는 재시도, 우선순위, 지연 실행 등은 Redis Streams 위에 구축된 추상화입니다.
직접 Stream을 다루지 않아도 되지만, BullMQ가 어떻게 동작하는지 이해하는 데 Stream 개념이 도움이 됩니다.
| 항목 | Pub/Sub | Stream | List |
|---|---|---|---|
| 메시지 저장 | 없음 | 영구 | 꺼내면 삭제 |
| 놓친 메시지 | 유실 | 재수신 가능 | 꺼낸 것만 유실 |
| 소비자 그룹 | 없음 | 있음 | 없음 |
| 처리 확인 | 없음 | ACK | 없음 |
| 적합한 용도 | 즉시 브로드캐스트 | 신뢰성 큐 | 간단한 큐 |
메시지 유실이 허용되면 Pub/Sub, 처리 보장이 필요하면 Stream, 단순 작업 큐면 List를 선택합니다.