RabbitMQ·Redis 큐 vs Kafka vs NATS: 메시징 시스템별 철학과 선택 기준
이 글의 핵심
결제 이벤트 한 건이 사라진 장애에서 출발해, 브로커가 말하는 exactly-once를 그대로 믿으면 안 되는 이유를 짚습니다. 작업 큐, 로그처럼 쌓는 Kafka, 가벼운 pub/sub인 NATS와 JetStream의 전달 보장 차이를 정리하고, 중복 키와 유니크 제약으로 소비 쪽을 멱등하게 만드는 실무 패턴을 제시합니다.
예전에 우리 팀이 배포한 날, 결제 이벤트가 한 건 사라진 적이 있었습니다. 대시보드에는 “프로듀서 성공”이 찍혀 있었지만 DB에는 기록이 없었고, 사용자는 돈을 냈는데 포인트만 쌓이지 않은 상태였습니다. 밤새 로그를 뒤져 보니 컨슈머가 메시지는 받았지만 (아직 ack 전에) OOM으로 죽었고, 그 사이 큐/브로커 쪽에서는 “전달 완료”로 보였는지, 반대로 “재시도”로 떠서 한 번 더 중복 전달된 것은 아닌지를 두고 논쟁이 벌어진 사례였습니다. 결론은 분명합니다. 딱 한 번(exactly-once)만 됩니다 같은 말은 브로커 UI에 떠 있어도 믿지 말아야 하고, 앱 레이어에서는 at-least-once가 현실입니다. 따라서 “중복이 올 수 있다”는 전제로 멱등 키와 부작용 처리를 설계해야 합니다.
RabbitMQ, Redis Queue, Kafka, NATS… 이름만 봐도 머리가 아픈 것은 누구나 마찬가지입니다. 다만 “작업을 나눠 처리한다”는 그림(메시지 큐), “이벤트를 로그처럼 쌓는다”는 그림(Kafka), “초고속 pub/sub”이라는 그림(NATS)은 짊어지는 부담이 서로 다릅니다. LinkedIn이 기존 메시징 시스템으로는 대량의 활동 로그를 감당하기 어려워 커밋 로그 = Kafka를 직접 만든 것도 그래서 유명한 이야기입니다. 느린 소비자가 브로커에 부담을 주느냐(Kafka는 pull + 로그), 아니면 큐에 밀려 쌓이다 터지느냐(전통 큐) 정도의 트레이드오프로 이해하면 됩니다.
실제로 exactly-once를 “시스템이 보장한다”고 쓰는 것은 홍보 문구에 가깝습니다. 프로듀서가 네트워크 끊김으로 재시도하면 컨슈머는 같은 메시지를 두 번 볼 수 있고, offset 커밋이 처리 뒤로 밀리면 장애 시 같은 메시지를 다시 읽습니다. 따라서 실무에서는 (1) at-least-once 전달, (2) 소비 쪽 멱등, (3) 외부 DB에 “이 주문_id 처리됨” 유니크 제약을 두는 방식이 정답에 가깝습니다. “한 번만” 처리해야 한다는 요구도 대부분 멱등 + 중복 키로 해결됩니다. exactly-once를 고집하다가 밤중에 호출 알림을 받는 것보다 훨씬 낫습니다.
RabbitMQ 쪽은 AMQP, Exchange, Binding, Queue… 튜토리얼대로 basic_ack를 잘 쓰더라도 큐가 쌓이면 메모리가 터집니다. Redis Queue는 LPUSH / BRPOP 정도로 단순하고 빠른 대신, 기본이 인메모리이고 영속성이 애매해서 “유실되면?”에 대한 대비를 직접 짜야 합니다. 두 경우 모두 질문은 비슷합니다. 소비가 끝나기 전에 프로세스가 죽으면 메시지가 다시 전달되는지, 다시 전달된다면 두 번 가는지입니다. 정답을 at-least-once + 멱등 쪽에 두는 팀이 오래 버팁니다.
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
channel = connection.channel()
channel.queue_declare(queue="tasks", durable=True)
channel.basic_publish(
exchange="",
routing_key="tasks",
body="Hello World",
properties=pika.BasicProperties(delivery_mode=2),
)
def callback(ch, method, properties, body):
print(f"Received {body}")
ch.basic_ack(delivery_tag=method.delivery_tag)
channel.basic_consume(queue="tasks", on_message_callback=callback)
channel.start_consuming()
이 코드에서 메시지를 “잃지 않기” 위한 장치는 세 군데에 나뉘어 있습니다. durable=True는 큐 자체가 브로커 재시작 뒤에도 남게 하고, delivery_mode=2는 메시지를 디스크에 쓰도록 표시하며, 콜백의 basic_ack는 처리가 끝났다는 사실을 브로커에 알립니다. 셋 중 하나만 빠져도 보장이 깨집니다. 큐만 durable이고 메시지가 persistent가 아니면 재시작 시 큐는 남지만 내용은 비어 있고, 반대로 큐가 durable이 아니면 persistent 메시지도 큐와 함께 사라집니다.
ack는 at-least-once의 핵심입니다. 컨슈머가 basic_ack를 보내기 전에 연결이 끊기면 RabbitMQ는 그 메시지를 다시 큐에 넣어 다른 컨슈머에게 보냅니다(이때 redelivered 플래그가 켜집니다). 서두의 장애처럼 “처리는 절반쯤 했는데 ack 전에 죽는” 경우 같은 메시지가 두 번 처리되는 이유가 이것입니다. 반대로 auto_ack=True로 소비하면 브로커가 보내는 순간 삭제하므로 at-most-once가 되어, 컨슈머가 죽으면 메시지가 그냥 사라집니다.
처음 이 구조를 운영할 때 흔히 겪는 문제는 basic_qos(prefetch_count=...)를 설정하지 않는 것입니다. 기본값에서는 브로커가 ack되지 않은 메시지를 컨슈머에게 제한 없이 밀어 넣기 때문에, 한 컨슈머가 수천 개를 메모리에 받아 둔 채 느리게 처리하고 새로 띄운 컨슈머는 놀게 됩니다. 그 상태에서 컨슈머가 OOM으로 죽으면 받아 둔 메시지 전체가 한꺼번에 재전달됩니다. 작업 큐라면 channel.basic_qos(prefetch_count=1)처럼 작게 시작해 처리량을 보며 늘리는 편이 안전합니다. 또 퍼블리셔 쪽에서 브로커가 메시지를 정말 받았는지 확인하려면 publisher confirms(channel.confirm_delivery())를 켜야 합니다. basic_publish가 예외 없이 끝났다는 것은 소켓에 썼다는 뜻일 뿐, 브로커가 디스크에 기록했다는 뜻이 아닙니다.
Redis를 큐로 쓸 때도 같은 질문에 답해야 합니다. BRPOP은 꺼내는 순간 리스트에서 원소가 사라지므로 처리 중에 워커가 죽으면 메시지가 유실됩니다. 이를 막는 고전적인 방법은 LMOVE(구 BRPOPLPUSH)로 꺼내면서 “처리 중” 리스트에 옮겨 두고 처리 후 지우는 것이고, 요즘은 컨슈머 그룹과 XACK, 미확인 메시지 목록(PEL)을 제공하는 Redis Streams를 쓰는 편이 낫습니다. 어느 쪽이든 Redis의 영속성은 AOF/RDB 설정에 달려 있어서, appendfsync everysec이면 장애 시 최대 1초 정도의 쓰기를 잃을 수 있다는 점을 전제로 설계해야 합니다.
Kafka는 토픽·파티션·offset으로 구성됩니다. 같은 user_id 키로 보내면 파티션 안에서는 순서가 보장됩니다. enable_auto_commit=True를 켜고 process 도중에 죽으면 “어디까지 처리했는지”와 “실제 부작용”이 어긋납니다. 그래서 운영을 잘하는 팀은 수동 커밋을 쓰거나 처리가 끝난 뒤 커밋하는 구조를 찾다가, 결국 “어차피 멱등”으로 돌아옵니다. at-least-once가 현실이라는 말이 여기서 한 번 더 나옵니다.
from kafka import KafkaConsumer
import json
consumer = KafkaConsumer(
"orders",
bootstrap_servers=["localhost:9092"],
auto_offset_reset="earliest",
enable_auto_commit=False, # 처리 후에 커밋 쪽이 안전에 가깝다
group_id="order-processor",
value_deserializer=lambda m: json.loads(m.decode("utf-8")),
)
for message in consumer:
try:
process_order(message.value)
consumer.commit()
except Exception as e:
# 주의: 여기서 커밋을 건너뛰어도 다음 메시지의 commit()이
# 이 메시지의 offset까지 함께 커밋한다 (자동 재시도 아님)
log_error(e)
enable_auto_commit=False로 두고 처리 후 commit()하는 구조는 “처리 전에 커밋되어 유실되는” 경우를 막아 줍니다. 자동 커밋은 일정 주기(auto_commit_interval_ms, 기본 5초)로 이미 가져온 위치를 커밋하므로, 가져온 메시지를 처리하는 도중 죽으면 처리하지 못한 메시지의 offset이 이미 커밋되어 있을 수 있습니다. 처리 후 커밋하면 반대로 “처리했는데 커밋 전에 죽어서 다시 읽는” 중복만 남고, 이것을 멱등성으로 흡수하는 것이 이 글의 기본 전략입니다.
그런데 이 예제의 except 블록은 흔히 오해하는 부분입니다. 주석에 “커밋 안 하면 재시도”라고 적는 경우가 많지만, Kafka의 offset은 메시지별 ack가 아니라 파티션별 위치 하나입니다. 실패한 메시지를 건너뛰고 루프를 계속 돌면 컨슈머의 위치는 이미 다음으로 넘어가 있고, 다음 메시지를 처리한 뒤의 commit()이 실패한 메시지까지 포함해 커밋합니다. 결국 실패한 메시지는 프로세스가 재시작되지 않는 한 다시 읽히지 않습니다. 실패를 정말 재시도하려면 consumer.seek(TopicPartition(message.topic, message.partition), message.offset)으로 위치를 되돌리거나, 몇 번 재시도한 뒤 별도의 재시도 토픽이나 DLQ 토픽으로 보내고 원래 흐름은 계속 진행하는 방식을 씁니다. 후자가 한 메시지 때문에 파티션 전체가 막히는(poison pill) 상황을 피할 수 있어 실무에서 더 흔합니다.
한 메시지마다 commit()을 동기로 호출하면 매번 브로커 왕복이 생겨 처리량이 크게 떨어진다는 점도 트레이드오프입니다. 보통은 몇백 개 단위 또는 몇 초 단위로 모아서 커밋하고, 그만큼 재시작 시 중복 처리되는 범위가 넓어지는 것을 멱등성으로 감당합니다. 또 처리 시간이 길어 max_poll_interval_ms(기본 5분)를 넘기면 브로커가 컨슈머를 죽은 것으로 보고 파티션을 다른 컨슈머에게 넘기는 리밸런스가 일어나며, 이때 커밋하지 못한 메시지들이 다른 컨슈머에서 다시 처리됩니다. “같은 주문이 두 번 처리됐는데 에러 로그는 없다”는 증상의 상당수가 이 리밸런스에서 나옵니다.
공정하게 말하면 Kafka의 exactly-once가 완전히 마케팅 문구인 것은 아닙니다. 멱등 프로듀서(enable.idempotence=true, Kafka 3.0부터 기본값)는 프로듀서 재시도로 인한 중복 기록을 브로커가 걸러 주고, 트랜잭션 API는 “토픽에서 읽고 → 처리하고 → 다른 토픽에 쓰고 → offset을 커밋하는” 과정을 원자적으로 묶어 줍니다. Kafka Streams의 processing.guarantee=exactly_once_v2가 이 위에서 동작합니다. 핵심은 이 보장의 범위가 Kafka 내부라는 점입니다. 처리 도중 외부 DB에 쓰거나 결제 API를 호출하는 순간 그 부작용은 Kafka 트랜잭션 밖에 있으므로, 결국 소비 쪽 멱등성이 다시 필요해집니다.
NATS는 가볍고 빠른 대신, 코어는 기본적으로 at-most-once에 가깝다는 점이 자주 지적됩니다(메모리, 재시작, 네트워크). 그래서 “절대 잃으면 안 되는 돈”이라면 JetStream을 쓰고 ack 모드를 확인하는 쪽을 봅니다. JetStream을 써도 앱이 같은 메시지를 이중으로 처리하는 것은 여전히 멱등 문제입니다. “JetStream = 영속”이라고 착각하면 앞의 결제 사고가 또 일어납니다.
import asyncio
from nats.aio.client import Client as NATS
async def main():
nc = NATS()
await nc.connect("nats://localhost:4222")
async def message_handler(msg):
print(f"{msg.subject}: {msg.data.decode()}")
await nc.subscribe("orders.*", cb=message_handler)
await nc.publish("orders.new", b"Order #123")
await asyncio.sleep(1)
await nc.close()
asyncio.run(main())
이 예제는 코어 NATS의 성격을 잘 보여 줍니다. subscribe가 먼저 호출되어 있으므로 같은 연결에서 발행한 orders.new를 받을 수 있지만, 구독자가 없는 시점에 발행된 메시지는 어디에도 저장되지 않고 버려집니다. 구독자가 잠깐 재연결 중이었거나 느려서 버퍼가 가득 찬 경우(slow consumer 에러)도 마찬가지입니다. orders.*의 *는 토큰 하나에 해당하는 와일드카드이고, orders.>는 그 아래 모든 계층을 받습니다. 마지막의 asyncio.sleep(1)은 콜백이 실행될 시간을 주려는 것으로, 실제 코드에서는 await nc.flush()나 await nc.drain()으로 보낸 메시지가 서버에 도달했는지 확인하고 연결을 정리하는 편이 정확합니다.
같은 과제를 JetStream으로 바꾸면 스트림이 메시지를 저장하고, 컨슈머가 명시적으로 ack해야 다음으로 넘어갑니다.
js = nc.jetstream()
await js.add_stream(name="ORDERS", subjects=["orders.*"])
# 발행: 서버가 저장했다는 응답(PubAck)을 기다린다
ack = await js.publish("orders.new", b"Order #123",
headers={"Nats-Msg-Id": "order-123"})
# 영속(durable) 컨슈머: 재시작해도 어디까지 처리했는지 서버가 기억
sub = await js.pull_subscribe("orders.*", durable="order-processor")
for msg in await sub.fetch(10, timeout=5):
process_order(msg.data)
await msg.ack()
Nats-Msg-Id 헤더는 JetStream의 중복 제거 기능으로, 설정된 시간 창(기본 2분) 안에 같은 ID로 다시 발행된 메시지를 스트림에 한 번만 저장합니다. 프로듀서 재시도로 인한 중복은 이것으로 줄일 수 있지만, 컨슈머가 처리 후 ack 전에 죽어 재전달되는 중복은 여전히 남습니다. ack 대기 시간(ack_wait, 기본 30초)을 처리 시간보다 짧게 두면 멀쩡히 처리 중인 메시지가 다시 전달되는 것도 흔한 설정 실수입니다.
소비 쪽을 멱등하게 만드는 방법
지금까지 계속 “멱등으로 막는다”고 했는데, 실제로는 메시지마다 고유한 키를 정하고, 처리 결과와 그 키를 같은 트랜잭션으로 기록하는 것이 가장 흔한 구현입니다. 키는 브로커의 메시지 ID보다 order_id나 payment_id처럼 비즈니스 이벤트를 식별하는 값이 좋습니다. 프로듀서가 재시도하면 브로커 입장에서는 서로 다른 메시지가 되어 ID가 달라질 수 있기 때문입니다.
import psycopg
def handle_payment_event(conn, event):
with conn.transaction():
cur = conn.execute(
"INSERT INTO processed_events (event_id) VALUES (%s) "
"ON CONFLICT (event_id) DO NOTHING",
(event["payment_id"],),
)
if cur.rowcount == 0:
return # 이미 처리한 이벤트: 조용히 건너뛴다
conn.execute(
"UPDATE users SET points = points + %s WHERE id = %s",
(event["points"], event["user_id"]),
)
# 트랜잭션이 커밋된 뒤에 브로커에 ack / offset 커밋
processed_events.event_id에 유니크 제약이 있으므로, 같은 이벤트가 두 번 와도 두 번째 INSERT는 아무 행도 넣지 못하고 포인트 적립도 건너뜁니다. 중요한 점은 “처리했다는 기록”과 “실제 부작용”이 같은 DB 트랜잭션에 들어 있다는 것입니다. 먼저 SELECT로 처리 여부를 확인하고 나중에 기록하는 방식은 두 컨슈머가 동시에 같은 메시지를 받으면 둘 다 “아직 처리 안 됨”을 보게 되어 경쟁 조건이 생깁니다. 유니크 제약에 판단을 맡기면 DB가 그 경쟁을 대신 해결해 줍니다.
부작용이 DB 밖에 있는 경우, 예를 들어 외부 결제 API나 이메일 발송이라면 이 방법만으로는 부족합니다. 외부 API가 멱등성 키(Idempotency-Key 헤더 같은)를 지원하면 이벤트 키를 그대로 넘기는 것이 가장 확실하고, 지원하지 않는다면 “보내기 전 상태 기록 → 발송 → 완료 기록”처럼 단계를 나누어 재시작 시 어느 단계부터 이어갈지 판단할 수 있게 만들어야 합니다. processed_events 테이블은 계속 커지므로, 브로커가 재전달할 수 있는 기간(Kafka 보존 기간, JetStream 중복 창 등)보다 충분히 긴 기간만 남기고 정리하는 작업도 필요합니다.
제가 이런 구조를 처음 도입할 때 가장 많이 놓치는 부분은 순서입니다. 컨슈머 코드에서 DB 커밋보다 ack(또는 offset 커밋)를 먼저 하면, 그 사이에 죽었을 때 “브로커는 처리 완료로 알고 DB에는 기록이 없는” 유실이 생깁니다. 서두의 결제 이벤트 장애가 정확히 이런 모양입니다. 규칙은 단순합니다. 부작용과 처리 기록을 먼저 커밋하고, ack는 가장 마지막에 보냅니다. 그러면 최악의 경우가 중복이 되고, 중복은 유니크 제약이 막아 줍니다.
선택 기준 정리
정리하자면, 벤치마크 숫자로 다툴수록 지치기만 합니다. 지연은 NATS < RabbitMQ·Kafka 정도의 그림이고, 처리량은 Kafka / NATS가 잘 나오는 편이지만, 그보다 “유실됐을 때 누가 곤란해지는가”를 먼저 정하는 것이 맞습니다. 백로그(큐 depth), 컨슈머 랙, DLQ, 파티션별 지연, 재처리율 같은 지표가 대시보드에 없다면 메시징을 쓰는 척만 하는 것입니다. 하이브리드(실시간은 NATS, 작업은 Rabbit, 로그/분석은 Kafka)도 괜찮습니다. 대신 각 시스템에서 기대하는 전달 시맨틱을 팀이 한 문장으로 말할 수 있어야 합니다. 저는 그 한 문장을 “최소 한 번은 옵니다. 중복은 감수하고 막는다”로 맞추는 편입니다.
같이 보면 이해가 빨라지는 내부 글: C++ 메시지 큐 · RabbitMQ·Kafka #50-7, Kafka 입문, C++ Kafka / librdkafka