RabbitMQ는 분산 시스템에서 메시지를 안정적으로 전달하는 데 널리 사용되는 강력한 메시지 브로커입니다. 하지만 때로는 예상치 못한 문제에 직면하여 개발자를 당황하게 만들기도 합니다. 특히 Python Pika 라이브러리를 사용하여 RabbitMQ를 연동할 때, 메시지가 의도치 않게 누락되거나 소비자가 모든 메시지를 처리하지 못하는 현상은 흔하지만 해결하기 까다로운 문제 중 하나입니다. 이 기술 블로그에서는 ‘소비자가 메시지의 절반만 받는’ 현상을 겪는 개발자를 위해, 문제의 원인을 분석하고 실질적인 해결 방안 및 예방 팁을 제공합니다.
1. 에러 발생 상황
이 문제는 RabbitMQ와 Python Pika를 사용하여 간단한 Producer-Consumer 패턴을 구현했을 때 발생했습니다. Producer가 총 20개의 메시지를 RabbitMQ 큐에 발행했지만, 유일하게 실행 중인 Consumer는 이 중 절반인 10개의 메시지(짝수 번호 메시지)만 수신했습니다. RabbitMQ 관리 콘솔에서는 모든 메시지가 큐에 정상적으로 존재하고 있는 것으로 확인되어, Producer 측의 문제는 아닌 것으로 보였습니다.
Producer 코드 예시
아래는 2초 간격으로 20개의 메시지를 발행하는 Python Producer 코드입니다. 메시지는 지속성을 위해 delivery_mode=2로 설정되었습니다.
import pika
import time
connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost', port=5672, credentials=pika.PlainCredentials(username='bert', password='bert')))
channel = connection.channel()
i = 0
while i < 20:
i = i + 1
message = f"XXX__{i}"
channel.basic_publish(exchange='', routing_key='testqueue', body=message, properties=pika.BasicProperties(delivery_mode=2))
time.sleep(2)
connection.close()
Consumer 코드 예시
다음은 메시지를 수신하고 출력하는 Python Consumer 코드입니다. 큐는 durable=True 및 x-queue-type: quorum으로 선언되었으며, auto_ack=True로 설정되어 메시지 수신 즉시 자동 승인됩니다.
import pika
import sys
import os
def callback(ch, method, properties, body):
print(" [x] Received %r" % body)
def main():
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost',
port=5672,
credentials=pika.PlainCredentials(username='bert',
password='bert')
)
)
channel = connection.channel()
channel.queue_declare(queue='testqueue', durable=True, arguments={'x-queue-type': 'quorum'})
channel.basic_consume(queue='testqueue', on_message_callback=callback, auto_ack=True)
print('Waiting for messages. ')
channel.start_consuming()
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
print('Interrupted')
try:
sys.exit(0)
except SystemExit:
os._exit(0)
Docker Compose 설정
RabbitMQ는 Docker Compose를 통해 실행되었습니다.
services:
rabbitmq:
image: rabbitmq:latest
container_name: rabbitmq
restart: always
ports:
- 5672:5672
- 15672:15672
environment:
RABBITMQ_DEFAULT_USER: bert
RABBITMQ_DEFAULT_PASS: bert
configs:
- source: rabbitmq-plugins
target: /etc/rabbitmq/enabled_plugins
volumes:
- ../data/rabbitmq/lib:/var/lib/rabbitmq/
- ../data/rabbitmq/logs:/var/log/rabbitmq
networks:
mynetwork:
ipv4_address: 10.0.1.13
configs:
rabbitmq-plugins:
content: "[rabbitmq_management]."
networks:
mynetwork:
external: true
name: mynet
Consumer 출력 결과
Producer 실행 후 Consumer에서 관찰된 출력은 다음과 같았습니다. 홀수 번호 메시지는 모두 누락되었습니다.
(.venv) $ py consume.py
Waiting for messages.
[x] Received b'XXX__2'
[x] Received b'XXX__4'
[x] Received b'XXX__6'
[x] Received b'XXX__8'
[x] Received b'XXX__10'
[x] Received b'XXX__12'
[x] Received b'XXX__14'
[x] Received b'XXX__16'
[x] Received b'XXX__18'
[x] Received b'XXX__20'
2. 명확한 발생 원인
제공된 코드 자체만으로는 소비자가 메시지의 절반만 받는 직접적인 논리적 오류를 찾기 어렵습니다. 특히 단일 소비자 환경에서 auto_ack=True는 메시지를 받은 즉시 승인하므로, 메시지가 큐에 남아있지만 소비자가 받지 못하는 상황은 흔치 않습니다. RabbitMQ 관리 콘솔에서도 모든 메시지가 존재한다고 보고되었으므로, Producer 측의 누락도 아니었습니다.
이러한 현상은 대부분 다음과 같은 환경적 또는 일시적인 상태 문제에서 발생할 가능성이 높습니다.
- 잔여(Stale) 또는 숨겨진 소비자(Phantom Consumer): 과거에 실행되었으나 정상적으로 종료되지 않았거나, 의도치 않게 백그라운드에서 실행되고 있는 다른 소비자 인스턴스가 존재하여 메시지를 경쟁적으로 소비했을 수 있습니다. 초기에는 단일 소비자라고 생각했더라도, 보이지 않는 곳에서 다른 소비자가 활성화되어 메시지를 가져갔을 가능성이 있습니다.
- 오래된 연결 또는 채널 상태: RabbitMQ 서버 또는 클라이언트 연결/채널이 불안정한 상태였을 수 있습니다. 때때로 연결이 끊어졌다가 재연결되는 과정에서 메시지 흐름이 일시적으로 불안정해지거나, 특정 메시지가 제대로 전달되지 않을 수 있습니다.
- RabbitMQ 서버의 일시적인 이상: 드물지만, RabbitMQ 브로커 자체의 일시적인 내부 상태 이상으로 인해 메시지 라우팅이나 전달에 문제가 발생했을 수 있습니다.
실제로 문제 해결 과정에서 “모든 것을 재시작”하고 모니터링했을 때 모든 메시지가 정상적으로 처리된 것을 보면, 코드 버그보다는 환경적인 요인 또는 일시적인 시스템 상태 오류였을 가능성이 매우 높습니다. rabbitmqctl list_consumers 명령을 통해 재시작 후 단 하나의 소비자만이 모든 메시지를 받았다는 사실이 이를 뒷받침합니다.
3. 해결 방법 및 코드 예시
이 문제의 핵심 해결책은 시스템 환경을 깨끗하게 재설정하는 것이었습니다. 이후에는 명시적인 메시지 승인(acknowledgment)을 사용하여 더욱 견고한 메시지 처리를 구현하는 것을 고려해볼 수 있습니다.
핵심 해결책: 시스템 완전 재시작 및 모니터링
가장 먼저 시도해야 할 방법은 Producer, Consumer, 그리고 RabbitMQ Docker 컨테이너를 포함한 모든 관련 구성 요소를 완전히 중지하고 재시작하는 것입니다. 이는 잔여 프로세스나 불안정한 연결 상태를 초기화하는 데 효과적입니다.
재시작 후, RabbitMQ 컨테이너 내부에서 rabbitmqctl 명령어를 사용하여 소비자의 상태를 직접 확인합니다. 이를 통해 현재 활성화된 소비자의 수와 어떤 큐를 소비하고 있는지 명확히 파악할 수 있습니다.
# RabbitMQ 컨테이너 내부에서 실행 (또는 Docker exec 사용)
rabbitmqctl list_consumers -p /
이 명령어를 통해 단 하나의 소비자가 testqueue를 소비하고 있으며, Producer가 보낸 20개의 메시지 모두를 성공적으로 수신하는 것을 확인했다면 문제는 해결된 것입니다. 이 경우, 별도의 코드 수정은 필요하지 않습니다.
보완적인 접근: 명시적 메시지 승인 (Manual Acknowledgment) 사용
만약 재시작 후에도 유사한 문제가 발생하거나, 더 견고한 메시지 처리 시스템을 구축하고자 한다면 auto_ack=False로 설정하고 메시지 처리 완료 후 명시적으로 승인하는 방식을 사용하는 것이 좋습니다. 이는 소비자가 메시지를 성공적으로 처리하기 전까지는 큐에서 메시지가 사라지지 않도록 보장하여 메시지 손실 위험을 최소화합니다.
Consumer 코드 수정 (auto_ack=False 적용)
import pika
import sys
import os
import time # time 모듈 추가
def callback(ch, method, properties, body):
print(" [x] Received %r" % body.decode()) # 바이트를 문자열로 디코딩
time.sleep(1) # 메시지 처리 시간 시뮬레이션
ch.basic_ack(method.delivery_tag) # 메시지 처리 완료 후 명시적으로 승인
def main():
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost',
port=5672,
credentials=pika.PlainCredentials(username='bert',
password='bert')
)
)
channel = connection.channel()
channel.queue_declare(queue='testqueue', durable=True, arguments={'x-queue-type': 'quorum'})
# auto_ack를 False로 설정하여 명시적 승인 사용
channel.basic_consume(queue='testqueue', on_message_callback=callback, auto_ack=False)
# prefetch_count 설정 (선택 사항): 소비자가 한 번에 가져갈 수 있는 메시지 수 제한
channel.basic_qos(prefetch_count=1)
print('Waiting for messages. ')
channel.start_consuming()
if __name__ == "__main__quot;:
try:
main()
except KeyboardInterrupt:
print('Interrupted')
try:
sys.exit(0)
except SystemExit:
os._exit(0)
위 코드에서는 auto_ack=False로 변경하고, callback 함수 내에서 ch.basic_ack(method.delivery_tag)를 호출하여 메시지를 명시적으로 승인하도록 했습니다. 또한, channel.basic_qos(prefetch_count=1)을 추가하여 소비자가 한 번에 하나의 메시지만 가져가도록 제한함으로써, 메시지가 공정하게 분배되고 과도한 프리패치로 인한 문제를 방지할 수 있습니다. 이는 여러 소비자가 있을 때 특히 중요하며, 단일 소비자 환경에서도 메시지 흐름을 보다 명확하게 제어하는 데 도움이 됩니다.
4. 향후 예방을 위한 팁
RabbitMQ와 같은 메시지 큐 시스템을 안정적으로 운영하기 위해서는 단순한 코드 구현을 넘어선 포괄적인 접근 방식이 필요합니다.
-
철저한 모니터링 습관화
RabbitMQ Management UI를 활용하여 큐의 상태, 메시지 수, 연결 상태, 소비자 목록 등을 주기적으로 확인하는 습관을 들여야 합니다. 또한,
rabbitmqctl명령어를 사용하여 서버의 현재 상태를 깊이 있게 들여다보는 것도 중요합니다. (예:rabbitmqctl list_queues name messages_ready consumers,rabbitmqctl list_connections). -
명시적 메시지 승인(Manual Acknowledgment) 사용
프로덕션 환경에서는
auto_ack=True의 편리함보다는auto_ack=False를 통한 명시적 메시지 승인 방식을 강력히 권장합니다. 메시지 처리 로직이 복잡하거나 외부 서비스와 연동될 때, 메시지가 완전히 처리되기 전까지는 큐에서 제거되지 않도록 하여 데이터 손실을 방지할 수 있습니다. -
안정적인 연결 관리 및 재연결 로직 구현
네트워크 불안정이나 RabbitMQ 서버 재시작 등 예측 불가능한 상황에 대비하여 Producer와 Consumer 모두 견고한 재연결(reconnection) 로직을 갖추는 것이 중요합니다. Pika 라이브러리에는 비동기 연결(AsyncioConnection)을 통해 재연결을 보다 쉽게 구현할 수 있는 기능이 있습니다.
-
상세한 로깅(Logging)
Producer와 Consumer 양쪽에서 메시지 발행, 수신, 처리 과정에 대한 상세한 로그를 남기세요. 특히 오류가 발생했을 때 관련 메시지 ID, 시간, 발생한 문제 등을 기록하면 문제 해결에 결정적인 단서를 제공할 수 있습니다.
-
멱등성(Idempotency) 고려
메시지 큐 시스템에서는 메시지가 중복으로 처리될 가능성이 항상 존재합니다. 소비자가 동일한 메시지를 여러 번 처리해도 시스템에 부작용이 없도록(예: 데이터 중복 저장 방지, 결제 이중 처리 방지) 멱등성을 고려하여 애플리케이션을 설계해야 합니다.
이러한 모범 사례들을 통해 RabbitMQ를 활용한 메시징 시스템을 더욱 안정적이고 효율적으로 운영할 수 있을 것입니다. 이번 문제를 계기로 더욱 견고한 시스템을 구축하는 기회가 되시길 바랍니다.
![[에러 해결] [Python asyncio] Redis PUB/SUB 장기 실행 TaskGroup 사용의 함정과 해결책 [에러 해결] [Python asyncio] Redis PUB/SUB 장기 실행 TaskGroup 사용의 함정과 해결책](https://dev-error.com/wp-content/plugins/contextual-related-posts/default.png)