RocketMQ 네이티브 API 사용법

일반 메시지

생산자

단방향 전송

메시지를 한 방향으로만 전송하며, 전송 성공 여부를 확인할 수 없습니다.

// 생산자 그룹 지정하여 생산자 생성
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;
    }
}

태그: RocketMQ java Message Queue Producer Consumer

8월 3일 19:36에 게시됨