메시지 큐(MQ) 핵심 구성 요소 및 구현

Docker 설치

docker run -d \
  --name rabbitmq-server \
  --hostname rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  -e RABBITMQ_DEFAULT_USER=admin \
  -e RABBITMQ_DEFAULT_PASS=admin123 \
  -v mq-data:/var/lib/rabbitmq \
  --network rabbitmq-net \
  rabbitmq:3.8-management

관리 콘솔(15672) 및 메시지 전송 인터페이스(5672) 포트를 노출하고, 데이터 볼륨과 전용 네트워크를 구성합니다.

아키텍처 핵심 요소

  • 생산자(Producer): 메시지 생성 및 전송 역할
  • 소비자(Consumer): 메시지 처리 및 소비 역할
  • 큐(Queue): 메시지 임시 저장 공간
  • 교환기(Exchange): 라우팅 결정 및 메시지 분배
  • 가상 호스트(Virtual Host): 논리적 격리 단위

SpringAMQP 활용

AMQP 프로토콜 기반으로 다국어 지원이 가능하나, Java 클라이언트의 복잡성을 해결하기 위해 SpringAMQP을 적용합니다. 주요 기능은 다음과 같습니다:

  • 큐/교환기 자동 생성 및 바인딩 관리
  • @RabbitListener 애노테이션을 통한 비동기 메시지 수신
  • RabbitTemplate을 통한 간편한 메시지 전송
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

워크 큐 모델

동일한 큐에 여러 소비자를 바인딩하여 메시지 처리 부하를 분산합니다. 처리 속도 불균형 문제를 해결하기 위해 프리페치 설정이 필수적입니다.

연결 설정

spring:
  rabbitmq:
    host: 10.0.0.5
    port: 5672
    virtual-host: /task-system
    username: task-user
    password: securePass!@#

메시지 발송 예제

@Test
public void publishTaskMessages() {
    String queue = "task-queue";
    for (int i = 0; i < 100; i++) {
        rabbitTemplate.convertAndSend(queue, "task-" + i);
        Thread.sleep(15);
    }
}

소비자 구현

@RabbitListener(queues = "task-queue")
public void handleFastTask(String payload) {
    System.out.println("빠른 처리: " + payload + " @ " + LocalTime.now());
    Thread.sleep(10);
}

@RabbitListener(queues = "task-queue")
public void handleSlowTask(String payload) {
    System.err.println("느린 처리: " + payload + " @ " + LocalTime.now());
    Thread.sleep(500);
}

프리페치 최적화

spring:
  rabbitmq:
    listener:
      simple:
        prefetch: 1

교환기 유형 비교

  • Fanout: 모든 바인딩 큐에 브로드캐스트
  • Direct: 정확한 라우팅 키 매칭
  • Topic: 와일드카드(#{0+}, *{1})를 이용한 유연한 매칭
  • Headers: 메시지 헤더 필드 기반 매칭 (특수 용도)

Fanout 교환기 설정

@Configuration
public class FanoutConfig {
    @Bean
    public FanoutExchange fanoutExchange() {
        return new FanoutExchange("app.fanout");
    }

    @Bean
    public Queue taskQueue() {
        return new Queue("queue.fanout.1");
    }

    @Bean
    public Binding fanoutBinding(Queue queue, FanoutExchange exchange) {
        return BindingBuilder.bind(queue).to(exchange);
    }
}

Topic 교환기 애노테이션 설정

@RabbitListener(
    bindings = @QueueBinding(
        value = @Queue(name = "queue.topic.1"),
        exchange = @Exchange(name = "app.topic", type = ExchangeTypes.TOPIC),
        key = "asia.*"
    )
)
public void processTopicMessage(String msg) {
    System.out.println("토픽 메시지: " + msg);
}

JSON 직렬화 최적화

JDK 직렬화 대신 Jackson을 사용하여 데이터 크기 감소 및 안전성 향상

@Bean
public MessageConverter jsonConverter() {
    Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
    converter.setCreateMessageIds(true);
    return converter;
}

태그: RabbitMQ spring-amqp work-queues exchange-types JSON-Serialization

9월 28일 05:59에 게시됨