이 문서에서는 Spring Boot 애플리케이션에서 Apache RocketMQ를 사용하여 간단한 메시지 생산자(Producer)와 소비자(Consumer)를 설정하는 방법을 단계별로 설명합니다. RocketMQ Spring Boot Starter를 활용하여 프로젝트를 구성합니다.
1. Maven 의존성 추가
먼저 pom.xml에 RocketMQ Spring Boot Starter 의존성을 추가합니다. 이를 통해 자동 설정과 편리한 API를 사용할 수 있습니다.
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version>
</dependency>
2. 메시지 생산자(Producer) 구현
생산자는 RocketMQTemplate을 주입받아 다양한 방식(동기, 비동기, 순서 보장)으로 메시지를 전송할 수 있습니다. 아래 예제는 비동기 및 순서 보장 전송을 구현한 클래스입니다.
import com.alibaba.fastjson.JSON;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class MessageProducer {
private static final Logger log = LoggerFactory.getLogger(MessageProducer.class);
@Resource
private RocketMQTemplate rocketMQTemplate;
/**
* 비동기 메시지 전송
*/
public void sendAsync(String topic, Object payload) {
rocketMQTemplate.asyncSend(topic, payload, new SendCallback() {
@Override
public void onSuccess(SendResult result) {
log.info("비동기 전송 성공: payload = {}, 상태 = {}",
JSON.toJSONString(payload), result.getSendStatus());
}
@Override
public void onException(Throwable ex) {
log.error("비동기 전송 실패: {}", ex.getMessage());
}
});
}
/**
* 순서 보장 비동기 메시지 전송
* 동일한 hashKey를 가진 메시지는 동일한 큐로 전송되어 순서가 유지됩니다.
*/
public void sendAsyncOrdered(String topic, Object payload, String hashKey) {
rocketMQTemplate.asyncSendOrderly(topic, payload, hashKey, new SendCallback() {
@Override
public void onSuccess(SendResult result) {
log.info("순서 보장 전송 성공: payload = {}, 상태 = {}",
JSON.toJSONString(payload), result.getSendStatus());
}
@Override
public void onException(Throwable ex) {
log.error("순서 보장 전송 실패: {}", ex.getMessage());
}
});
}
}
3. 메시지 소비자(Consumer) 구현
소비자는 @RocketMQMessageListener 어노테이션을 사용하여 특정 토픽과 소비자 그룹을 구독합니다. 아래 예제는 순서 보장 모드로 메시지를 처리하는 기본 소비자입니다.
import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
@RocketMQMessageListener(
topic = "${custom.topic.name:default-topic}",
consumerGroup = "${custom.consumer.group:default-group}",
consumeMode = ConsumeMode.ORDERLY
)
@Component
public class MessageConsumer implements RocketMQListener<String> {
private static final Logger log = LoggerFactory.getLogger(MessageConsumer.class);
@Override
public void onMessage(String message) {
log.info("수신 메시지: {}, 비즈니스 로직 처리 중...", message);
// 여기서 실제 비즈니스 로직을 구현합니다.
}
}
4. 애플리케이션 설정 (application.yml)
RocketMQ 관련 연결 정보와 동작 파라미터를 application.yml에 정의합니다. 실제 환경에 맞게 NameServer 주소와 그룹명을 수정하십시오.
rocketmq:
# NameServer 주소 (클러스터인 경우 세미콜론으로 구분)
name-server: 192.168.0.1:9876;192.168.0.2:9876
# 생산자(Producer) 설정
producer:
# 생산자 그룹
group: my-producer-group
# 메시지 전송 타임아웃 (ms, 기본값 3000)
send-message-timeout: 3000
# 메시지 압축 임계값 (bytes, 기본값 4096)
compress-message-body-threshold: 4096
# 최대 메시지 크기 (bytes, 기본값 4194304)
max-message-size: 4194304
# 동기 전송 실패 시 재시도 횟수 (기본값 2)
retry-times-when-send-failed: 2
# 비동기 전송 실패 시 재시도 횟수 (기본값 2)
retry-times-when-send-async-failed: 2
# 실패 시 다른 Broker 재시도 여부 (기본값 false)
retry-next-server: false
# 메시지 추적 활성화 여부 (기본값 true)
enable-msg-trace: false
# 사용자 정의 추적 토픽 (기본값 RMQ_SYS_TRACE_TOPIC)
customized-trace-topic: RMQ_SYS_TRACE_TOPIC
# 사용자 정의 설정: 토픽 및 소비자 그룹
my-topic: order-topic
my-consumer-group: order-group
5. 테스트 컨트롤러를 통한 메시지 전송 확인
HTTP 요청을 통해 메시지를 전송하고, 소비자 로그를 확인하는 간단한 컨트롤러를 작성합니다.
import org.springframework.beans.factory.annotation.Value;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
@RestController
public class MessageTestController {
@Resource
private MessageProducer producer;
@Value("${rocketmq.my-topic:default-topic}")
private String topic;
@GetMapping("/sendOrderedMessage")
public String sendMessage(@RequestParam String content) {
// hashKey는 메시지의 고유 식별자를 사용하여 동일한 큐로 전송되도록 합니다.
// 예: 주문 ID 또는 사용자 ID
String hashKey = "fixed-hash-key";
producer.sendAsyncOrdered(topic, content, hashKey);
return "메시지 전송 완료: " + content;
}
}
애플리케이션을 실행하고 GET /sendOrderedMessage?content=test를 호출하면, 생산자가 비동기로 메시지를 전송하고 소비자 로그에서 수신 여부를 확인할 수 있습니다.