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;
}