일반 메시지
생산자
단방향 전송
메시지를 한 방향으로만 전송하며, 전송 성공 여부를 확인할 수 없습니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("oneway-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
// 토픽 지정하여 메시지 생성
Message message = new Message("onewayTopic", "단방향 메시지입니다.".getBytes());
// 단방향 전송 (성공 여부 반환 없음)
producer.sendOneway(message);
// 생산자 종료
producer.shutdown();
동기 전송
메시지를 전송하고 서버의 응답을 기다립니다. 전송 결과를 반환받을 수 있습니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("test-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
// 10개의 메시지 전송 (4개의 큐에 분산됨: 2-2-3-3)
for (int i = 0; i < 10; i++) {
// 토픽 지정 및 메시지 생성
Message message = new Message("testTopic", "간단한 메시지입니다.".getBytes());
// 동기 전송 (메시지 ID, 전송 큐, 오프셋 등의 정보 포함된 SendResult 반환)
SendResult sendResult = producer.send(message);
System.out.println("메시지 전송 상태: " + sendResult.getSendStatus());
}
// 생산자 종료
producer.shutdown();
비동기 전송
메시지를 전송하고 즉시 반환되며, 서버 응답은 콜백 함수를 통해 처리됩니다. 현재 스레드를 차단하지 않습니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("async-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
// 토픽 지정하여 메시지 생성
Message message = new Message("asyncTopic", "비동기 메시지입니다.".getBytes());
// 비동기 전송 (SendCallback을 통해 서버 응답 처리)
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
System.out.println("전송 성공");
}
@Override
public void onException(Throwable e) {
System.err.println("전송 실패: " + e.getMessage());
}
});
// 생산자 종료
producer.shutdown();
소비자
메시지를 구독하고 처리하는 컴포넌트입니다.
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("test-consumer-group");
// 네임서버 주소 설정
consumer.setNamesrvAddr("127.0.0.1:9876");
// 토픽 구독 ( '*'는 모든 태그를 의미, 태그 필터링 가능)
consumer.subscribe("testTopic", "*");
// 메시지 리스너 등록 (MessageListenerConcurrently: 동시성 모드)
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
// msgs 리스트에는 여러 메시지가 포함될 수 있습니다.
MessageExt msg = msgs.get(0); // 첫 번째 메시지 가져오기
System.out.println("메시지 전체 내용: " + msg.toString());
// 메시지 본문 디코딩
System.out.println("메시지 본문: " + new String(msg.getBody()));
System.out.println("소비 컨텍스트: " + context);
// CONSUME_SUCCESS 반환 시 메시지는 MQ에서 제거됩니다.
// RECONSUME_LATER 반환 시 메시지는 재소비를 위해 큐로 반환됩니다 (기본 16회 재시도).
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 소비자 시작
consumer.start();
배치 메시지
여러 메시지를 묶어 한 번에 전송하지만, 소비자는 여전히 메시지를 개별적으로 처리합니다.
생산자
여러 메시지를 리스트로 묶어 한 번의 API 호출로 전송합니다. 이 메시지들은 동일한 큐로 전송됩니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("batch-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
// 전송할 메시지 리스트 생성
List<Message> messages = Arrays.asList(
new Message("batchTopic", "배치 메시지 A".getBytes()),
new Message("batchTopic", "배치 메시지 B".getBytes()),
new Message("batchTopic", "배치 메시지 C".getBytes())
);
// 메시지 리스트를 전송 (단일 API 호출)
SendResult sendResult = producer.send(messages);
System.out.println("배치 메시지 전송 결과: " + sendResult);
// 생산자 종료
producer.shutdown();
소비자
배치로 전송된 메시지를 소비할 때, `consumeMessage` 메서드에 여러 메시지가 리스트 형태로 전달될 수 있습니다.
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("batch-consumer-group");
// 네임서버 주소 설정
consumer.setNamesrvAddr("127.0.0.1:9876");
// 토픽 구독
consumer.subscribe("batchTopic", "*");
// 메시지 리스너 등록
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
// msgs 리스트에는 여러 메시지가 포함될 수 있습니다.
System.out.println("메시지 수신 시간: " + new Date());
System.out.println("수신된 메시지 개수: " + msgs.size());
// 리스트의 첫 번째 메시지 본문 출력
System.out.println("첫 번째 메시지 본문: " + new String(msgs.get(0).getBody()));
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 소비자 시작
consumer.start();
지연 메시지
메시지가 즉시 소비되지 않고 지정된 시간 후에 소비되도록 설정할 수 있습니다.
- RocketMQ 4.x 버전에서는 미리 정의된 지연 레벨을 사용합니다. 이 레벨은 설정 파일에서 수정 가능합니다.
- RocketMQ 5.x 버전에서는 특정 시간을 지정할 수 있습니다.
- 소비자 측에서는 일반 메시지와 동일하게 처리합니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("delay-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
// 토픽 지정 및 지연 메시지 생성
Message message = new Message("delayTopic", "지연 메시지입니다.".getBytes());
// 지연 시간 레벨 설정 (예: 레벨 3은 설정된 지연 시간 중 세 번째 시간을 의미)
// messageDelayLevel 설정 예: "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"
message.setDelayTimeLevel(3);
// 메시지 전송
producer.send(message);
System.out.println("지연 메시지 전송 시간: " + new Date());
// 생산자 종료
producer.shutdown();
순서 메시지
메시지가 전송된 순서대로 소비되어야 할 때 사용합니다. 이를 위해 생산자는 메시지를 동일한 큐로 보내고, 소비자는 단일 스레드 모드로 메시지를 처리해야 합니다.
생산자
동일한 비즈니스 키(예: 주문 ID)를 가진 메시지들이 항상 같은 메시지 큐로 전송되도록 `MessageQueueSelector`를 구현합니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("ordered-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
String[] orderIds = {"ORDER_001", "ORDER_002", "ORDER_003"};
String[] orderSteps = {"CREATE", "PAY", "SHIP", "COMPLETE"};
for (String orderId : orderIds) { // 각 주문에 대해
for (String step : orderSteps) { // 각 단계를 순서대로
// 토픽, 태그, 메시지 본문 설정
Message msg = new Message("OrderTopic", "Order", (orderId + ":" + step).getBytes());
// 순서 메시지 전송 (MessageQueueSelector 사용)
producer.send(msg, new MessageQueueSelector() {
// 어떤 메시지 큐로 보낼지 결정하는 로직
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
String orderId = (String) arg; // 전달된 주문 ID
// 주문 ID의 해시코드를 사용하여 큐 인덱스 결정
int index = Math.abs(orderId.hashCode()) % mqs.size();
return mqs.get(index);
}
}, orderId); // 주문 ID를 select 메서드의 인자로 전달
System.out.printf("순서 메시지 전송: %s%n", new String(msg.getBody()));
Thread.sleep(500); // 처리 간격 시뮬레이션
}
}
// 생산자 종료
producer.shutdown();
소비자
순서 메시지를 처리하기 위해 `MessageListenerOrderly`를 사용합니다. 이 모드는 단일 스레드로 메시지를 소비하여 순서를 보장합니다.
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ordered-consumer-group");
// 네임서버 주소 설정
consumer.setNamesrvAddr("localhost:9876");
// 토픽 및 태그 구독
consumer.subscribe("OrderTopic", "Order");
// MessageListenerOrderly: 순서 모드 (단일 스레드, 무한 재시도 기본값)
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
for (MessageExt msg : msgs) {
String messageBody = new String(msg.getBody());
System.out.printf("처리 중인 메시지: %s%n", messageBody);
// 실제 비즈니스 로직 처리
}
// SUCCESS 반환 시 해당 메시지는 처리 완료됩니다.
return ConsumeOrderlyStatus.SUCCESS;
}
});
// 소비자 시작
consumer.start();
System.out.println("순서 소비자 시작됨.");
메시지 태그 필터링
메시지 전송 시 태그를 부여하고, 소비자는 특정 태그의 메시지만 구독하여 필터링할 수 있습니다.
생산자
다른 태그를 가진 두 개의 메시지를 전송합니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("tag-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
// 태그 'vip1'을 가진 메시지 생성 및 전송
Message messageVip1 = new Message("tagTopic", "vip1", "VIP 1 메시지".getBytes());
producer.send(messageVip1);
// 태그 'vip2'을 가진 메시지 생성 및 전송
Message messageVip2 = new Message("tagTopic", "vip2", "VIP 2 메시지".getBytes());
producer.send(messageVip2);
// 생산자 종료
producer.shutdown();
소비자 1
태그가 'vip1'인 메시지만 구독합니다.
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer1 = new DefaultMQPushConsumer("tag-consumer-group-a");
// 네임서버 주소 설정
consumer1.setNamesrvAddr("127.0.0.1:9876");
// 태그 'vip1'만 구독
consumer1.subscribe("tagTopic", "vip1");
// 메시지 리스너 등록
consumer1.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
System.out.println("소비자 1: " + new String(msgs.get(0).getBody()));
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 소비자 시작
consumer1.start();
소비자 2
'vip1' 또는 'vip2' 태그를 가진 모든 메시지를 구독합니다.
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer2 = new DefaultMQPushConsumer("tag-consumer-group-b");
// 네임서버 주소 설정
consumer2.setNamesrvAddr("127.0.0.1:9876");
// 태그 'vip1' 또는 'vip2' 구독 (OR 연산자 사용)
consumer2.subscribe("tagTopic", "vip1 || vip2");
// 메시지 리스너 등록
consumer2.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
System.out.println("소비자 2: " + new String(msgs.get(0).getBody()));
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 소비자 시작
consumer2.start();
메시지 Key 지정
메시지에 고유한 키를 부여하여 나중에 메시지를 조회하거나 추적하는 데 사용할 수 있습니다.
생산자
메시지와 함께 고유한 Key를 전송합니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("key-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
String uniqueKey = UUID.randomUUID().toString();
// 토픽, 태그, Key, 메시지 본문 설정
Message message = new Message("keyTopic", "vip1", uniqueKey, "Key가 포함된 메시지".getBytes());
producer.send(message);
// 생산자 종료
producer.shutdown();
소비자
수신된 메시지에서 Key를 추출하여 확인할 수 있습니다.
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("key-consumer-group");
// 네임서버 주소 설정
consumer.setNamesrvAddr("127.0.0.1:9876");
// 토픽 구독
consumer.subscribe("keyTopic", "*");
// 메시지 리스너 등록
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
MessageExt messageExt = msgs.get(0);
System.out.println("메시지 본문: " + new String(messageExt.getBody()));
System.out.println("메시지 Key: " + messageExt.getKeys());
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 소비자 시작
consumer.start();
재시도 메커니즘
메시지 전송 또는 소비 실패 시 자동으로 재시도하는 기능입니다.
생산자 재시도
메시지 전송 실패 시 재시도 횟수를 설정합니다.
// 생산자 그룹 지정하여 생산자 생성
DefaultMQProducer producer = new DefaultMQProducer("retry-producer-group");
// 네임서버 주소 설정
producer.setNamesrvAddr("127.0.0.1:9876");
// 생산자 시작
producer.start();
// 동기 전송 실패 시 재시도 횟수 설정
producer.setRetryTimesWhenSendFailed(2);
// 비동기 전송 실패 시 재시도 횟수 설정
producer.setRetryTimesWhenSendAsyncFailed(2);
Message message = new Message("retryTopic", "vip1", "재시도 메시지입니다.".getBytes());
producer.send(message);
// 생산자 종료
producer.shutdown();
소비자 재시도
동시성 모드 (MessageListenerConcurrently)
- 최대 재시도 횟수: 16회 (기본값)
- 재시도 간격: 10초, 30초, 1분, ..., 2시간 (지연 레벨에 따라 다름)
- 소비 실패 시, 최대 재시도 횟수를 초과하지 않으면 메시지는 재시도 큐(`%RETRY%<소비자 그룹 이름>`)로 이동합니다. `getReconsumeTimes()` 메서드로 현재 재시도 횟수를 얻을 수 있습니다.
- 최대 재시도 횟수를 초과하면 메시지는 데드 레터 큐(`%DLQ%<소비자 그룹 이름>`)로 이동합니다.
- 데드 레터 큐의 메시지를 처리하기 위해 해당 토픽을 구독하는 별도의 소비자를 생성해야 합니다.
- 비즈니스 로직에서 특정 재시도 횟수 후 수동 개입을 위해 로깅할 수 있습니다.
/**
* 데드 레터 큐의 토픽: %DLQ%retry-consumer-group
*/
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("retry-consumer-group");
// 네임서버 주소 설정
consumer.setNamesrvAddr("127.0.0.1:9876");
// 토픽 구독
consumer.subscribe("retryTopic", "*");
// 동시성 모드: 최대 재시도 횟수 설정
consumer.setMaxReconsumeTimes(10);
// 동시성 모드: 다음 소비 시 재시도 큐의 지연 레벨 설정
consumer.setDelayLevelWhenNextConsume(3);
// 메시지 리스너 등록
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
MessageExt messageExt = msgs.get(0);
// 현재 메시지의 재시도 횟수 (동시성 모드에서만 유효)
int retryCount = messageExt.getReconsumeTimes();
if (retryCount > 3) { // 3회 실패 시 성공으로 간주하고 처리 중단 (로깅 필요)
System.out.println("3회 이상 재시도 실패, 성공 처리: " + messageExt.getMsgId());
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
try {
// 실제 비즈니스 로직 처리
System.out.println("메시지 처리 시도, 재시도 횟수: " + retryCount);
// ... 로직 ...
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
// 비즈니스 로직 처리 중 오류 발생 시, 메시지는 재시도 큐로 이동
System.err.println("메시지 처리 오류: " + messageExt.getMsgId() + ", 재시도 횟수: " + retryCount);
return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 재소비 요청
}
}
});
// 소비자 시작
consumer.start();
순서 모드 (MessageListenerOrderly)
- 최대 재시도 횟수: 무제한 (Integer.MAX_VALUE)
- 재시도 간격: 기본 1000ms (1초)
- 동시성 모드와 달리, 재시도 시 별도의 재시도 큐로 이동하지 않고 원래 큐에서 재시도합니다. 이로 인해 해당 큐의 메시지 처리가 차단될 수 있습니다.
- 데드 레터 큐로 이동하지 않습니다.
// 각 큐별 오류 카운터를 관리하는 맵
private ConcurrentMap<Integer, AtomicInteger> queueErrorCounters = new ConcurrentHashMap<>();
private final int MAX_RETRY = 5; // 최대 허용 재시도 횟수 (예시)
// 소비자 그룹 지정하여 소비자 생성
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ordered-retry-consumer-group");
// 네임서버 주소 설정
consumer.setNamesrvAddr("127.0.0.1:9876");
// 토픽 구독
consumer.subscribe("retryTopic", "*");
// 순서 모드: 현재 큐 일시 중단 시간 설정 (밀리초)
consumer.setSuspendCurrentQueueTimeMillis(5000);
// 메시지 리스너 등록 (순서 모드)
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
return consumeMessagesSequentially(msgs, context);
}
});
// 소비자 시작
consumer.start();
// 실제 메시지 처리 로직
private ConsumeOrderlyStatus consumeMessagesSequentially(List<MessageExt> msgs, ConsumeOrderlyContext context) {
MessageExt msg = msgs.get(0);
int queueId = msg.getQueueId(); // 메시지가 속한 큐 ID
try {
// 실제 비즈니스 로직 처리
System.out.println("순서 메시지 처리 중: " + new String(msg.getBody()));
// 성공 시 해당 큐의 오류 카운터 제거
queueErrorCounters.remove(queueId);
return ConsumeOrderlyStatus.SUCCESS;
} catch (Exception e) {
// 오류 발생 시
// 해당 큐의 오류 카운터 증가
AtomicInteger counter = queueErrorCounters.computeIfAbsent(queueId, k -> new AtomicInteger(0));
int retryCount = counter.incrementAndGet();
// 최대 재시도 횟수 초과 시
if (retryCount > MAX_RETRY) {
System.err.println("메시지 [" + msg.getMsgId() + "] 재시도 횟수 초과 (" + MAX_RETRY + "회), 수동 개입 필요.");
// 주의: SUCCESS를 반환하면 이 메시지는 다시 소비되지 않으며, 데이터 불일치가 발생할 수 있습니다.
// 실제 환경에서는 로깅 또는 알림 후 SUCCESS를 반환하거나 다른 처리를 해야 합니다.
return ConsumeOrderlyStatus.SUCCESS; // 임시 방편
}
// 현재 큐 일시 중단 요청
System.out.println("메시지 [" + msg.getMsgId() + "] 처리 실패, " + retryCount + "회 재시도. 큐 일시 중단.");
return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
}
}