Spring Boot集成RocketMQ生产者消费者

PunkLu 2019年12月09日 66次浏览
Spring Boot集成RocketMQ生产者消费者
# 创建Spring Boot项目

使用IDEA自带的创建Spring Boot功能:

File -> New -> Project -> Spring Initializr

创建完成后,添加依赖:

<!-- https://mvnrepository.com/artifact/org.apache.rocketmq/rocketmq-client -->
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>4.5.2</version>
</dependency>

创建目录结构

  1. 创建constants包
  2. 创建quickstart包

创建生产者

创建生产者类

​ 在quickstart包下新建Producer类:

import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.exception.RemotingException;
import tech.punklu.rocketmqapi.constants.Const;

/**
 * 消息生产者
 */
public class Producer {

    public static void main(String[] args) throws MQClientException, RemotingException, InterruptedException, MQBrokerException {
        // 创建默认MQ消息生产者,test_quick_producer_name为生产者组的名称
        DefaultMQProducer producer = new DefaultMQProducer("test_quick_producer_name");
        // 绑定NameServer地址
        producer.setNamesrvAddr(Const.NAMESRV_ADDR);
        // 启动
        producer.start();
        for (int i = 0; i < 5; i++) {
            // 创建消息
            Message message = new Message("test_quick_topic",   // 主题
                    "TagA", // 标签
                    "keyA" + i, // 用户自定义的key,唯一的标识
                    ("Hello RocketMQ" + i).getBytes()); // 消息内容实体(byte[])
            // 发送消息
            SendResult sendResult  = producer.send(message);
            System.out.println("消息发出" + sendResult);
        }

        // 生产者关闭
        producer.shutdown();
    }
}

创建消费者类

​ 在quickstart目录下新建Consumer类:

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import tech.punklu.rocketmqapi.constants.Const;

import java.util.List;

/**
 * 消息消费者
 */
public class Consumer {

    public static void main(String[] args) throws MQClientException {
        // 创建默认消费者组,test_quick_consumer_name为对应的消费者组的名称
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("test_quick_consumer_name");
        // 设置NameServer
        consumer.setNamesrvAddr(Const.NAMESRV_ADDR);
        // 设置从最后端开始消费
        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
        // 订阅,*表示订阅此topic下的所有消息
        consumer.subscribe("test_quick_topic","*");

        // 注册消息消费监听器,内部实现的consumeMessage是对消息的消费方法
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list,
                                                            ConsumeConcurrentlyContext consumeConcurrentlyContext) {
                MessageExt messageExt = list.get(0);
                try {
                    String topic = messageExt.getTopic();
                    String tags = messageExt.getTags();
                    String keys = messageExt.getKeys();
                    if (keys.equals("keyA1")){
                        System.out.println("消息消费失败...");
                        int a = 1/0;
                    }
                    String msgBody = new String(messageExt.getBody(), RemotingHelper.DEFAULT_CHARSET);
                    System.out.println("topic: " + topic + ",tags: " + tags + ",keys " + keys + ",body: " + msgBody);
                }catch (Exception e){
                    e.printStackTrace();
                    // 获取当前消息已经重发的次数
                    int recousumeTimes = messageExt.getReconsumeTimes();
                    System.out.println("recousumeTimes: " + recousumeTimes);
                    if (recousumeTimes == 3){
                        // 记录日志...
                        // 做补偿处理
                        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                    }
                    // 稍后重试
                    return ConsumeConcurrentlyStatus.RECONSUME_LATER;
                }
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        consumer.start();
        System.out.println("consumer start....");
    }
}