RocketMQ的两种消费模式详解

 更新时间:2023年10月11日 11:07:35   作者:码奴生来只知道前进~  
这篇文章主要介绍了RocketMQ的两种消费模式详解,RocketMQ主要提供了两种消费模式,集群消费以及广播消费,我们只需要在定义消费者的时候通过setMessageModel(MessageModel.XXX),需要的朋友可以参考下

1、添加依赖

<dependency>
	<groupId>org.apache.rocketmq</groupId>
	<artifactId>rocketmq-client</artifactId>
	<version>4.4.0</version>
</dependency>
<dependency>
	<groupId>com.alibaba</groupId>
	<artifactId>fastjson</artifactId>
	<version>1.2.3</version>
</dependency>
<dependencies>
	<dependency>
		<groupId>org.projectlombok</groupId>
		<artifactId>lombok</artifactId>
	</dependency>
</dependencies>

2、消费模式

RocketMQ主要提供了两种消费模式:集群消费以及广播消费。我们只需要在定义消费者的时候通过setMessageModel(MessageModel.XXX)

// 设置消费模型,集群还是广播,默认为集群  CLUSTERING-集群,BROADCASTING-广播
mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);

方法就可以指定是集群还是广播式消费,默认是集群消费模式,即每个Consumer Group中的Consumer均摊所有的消息。

3、集群消费

3.1 生产者

package com.shucha.deveiface.biz.mq.producer;
import com.alibaba.fastjson.JSON;
import com.shucha.deveiface.biz.model.User;
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.common.RemotingHelper;
import java.util.Date;
/**
 * @author tqf
 * @Description 生产者
 * @Version 1.0
 * @since 2022-07-12 14:50
 */
public class MQProducer {
    public static void main(String[] args) throws MQClientException{
        producerSendMessage();
    }
    /**
     * 生产消息方法
     * @throws MQClientException
     */
    public static void producerSendMessage() throws MQClientException {
        // 创建DefaultMQProducer类并设定生产者名称
        DefaultMQProducer mqProducer = new DefaultMQProducer("producer-group-test");
        // 设置NameServer地址,如果是集群的话,使用分号;分隔开
        mqProducer.setNamesrvAddr("127.0.0.1:9876");
        // 消息最大长度 默认4M
        mqProducer.setMaxMessageSize(4096);
        // 发送消息超时时间,默认3000
        mqProducer.setSendMsgTimeout(3000);
        // 发送消息失败重试次数,默认2
        mqProducer.setRetryTimesWhenSendAsyncFailed(2);
        // 启动消息生产者
        mqProducer.start();
        try {
            // 循环十次,发送十条消息
            for (int i = 1; i <= 10; i++) {
                User user = new User();
                user.setId((long)i);
                user.setAge(i);
                user.setUserName("姓名"+i);
                user.setCreateTime(new Date());
                // String msg = "这是第" + i + "条消息测试";
                String msg = JSON.toJSONString(user);
                // 创建消息,并指定Topic(主题),Tag(标签)和消息内容
                Message message = new Message("TOPIC_TEST", "", msg.getBytes(RemotingHelper.DEFAULT_CHARSET));
                // 发送同步消息到一个Broker,可以通过sendResult返回消息是否成功送达
                SendResult sendResult = mqProducer.send(message);
                // mqProducer.sendOneway(message);
                // 消息id
                /*System.out.println(sendResult.getMsgId());
                // 队列信息
                System.out.println(sendResult.getMessageQueue());
                // 发送结果
                System.out.println(sendResult.getSendStatus());
                // 下一个要消费的消息的偏移量
                System.out.println(sendResult.getOffsetMsgId());
                // 队列消息偏移量
                System.out.println(sendResult.getQueueOffset());*/
                System.out.println(sendResult);
            }
        } catch (Exception e) {
            e.printStackTrace();
            System.out.println("生产消息异常!");
        }
        // 如果不再发送消息,关闭Producer实例
        mqProducer.shutdown();
    }
}

User用户测试实体类 

package com.shucha.deveiface.biz.model;
import com.fasterxml.jackson.annotation.JsonFormat;
import com.sdy.common.utils.DateUtil;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
import java.util.Date;
/**
 * @author tqf
 * @Description
 * @Version 1.0
 * @since 2022-04-07 13:53
 */
@Data
public class User {
    /**
     * 主键ID
     */
    private Long id;
    /**
     *用户名
     */
    private String userName;
    /**
     * 用户密码
     */
    private String passWord;
    /**
     * 年龄
     */
    private Integer age;
    /**
     * 性别(0-男,1-女,2-未知)
     */
    private Integer sex;
    /**
     * 创建时间
     */
    @JsonFormat(pattern = DateUtil.DATETIME_FORMAT)
    private Date createTime;
}

3.2 消费者A

package com.shucha.deveiface.biz.mq.consumer;
import com.alibaba.fastjson.JSON;
import com.shucha.deveiface.biz.constants.MqConstants;
import com.shucha.deveiface.biz.model.User;
import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
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.common.protocol.heartbeat.MessageModel;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import java.nio.charset.Charset;
import java.util.List;
/**
 * @author tqf
 * @Description 消费者A
 * @Version 1.0
 * @since 2022-07-12 14:37
 */
public class ConsumerA {
    public static void main(String[] args) throws MQClientException {
        ConsumerA();
    }
    public static void ConsumerA() throws MQClientException {
        // 创建DefaultMQPushConsumer类并设定消费者名称
        DefaultMQPushConsumer mqPushConsumer = new DefaultMQPushConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1);
        // DefaultMQPullConsumer pullConsumer = new DefaultMQPullConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1);
        // 设置NameServer地址,如果是集群的话,使用分号;分隔开
        mqPushConsumer.setNamesrvAddr("127.0.0.1:9876");
        // pullConsumer.setNamesrvAddr("127.0.0.1:9876");
        // 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费
        // 如果不是第一次启动,那么按照上次消费的位置继续消费
        mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
        // 设置消费模型,集群还是广播,默认为集群  CLUSTERING-集群   BROADCASTING-广播
        mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
        // mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
        // 消费者最小线程量
        mqPushConsumer.setConsumeThreadMin(5);
        // 消费者最大线程量
        mqPushConsumer.setConsumeThreadMax(10);
        // 设置一次消费消息的条数,默认是1
        mqPushConsumer.setConsumeMessageBatchMaxSize(1);
        // 订阅一个或者多个Topic,以及Tag来过滤需要消费的消息,如果订阅该主题下的所有tag,则使用*
        mqPushConsumer.subscribe("TOPIC_TEST", "*");
        // 注册回调实现类来处理从broker拉取回来的消息
        mqPushConsumer.registerMessageListener(new MessageListenerConcurrently() {
            // 监听类实现MessageListenerConcurrently接口即可,重写consumeMessage方法接收数据
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgList, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
                for (MessageExt msg : msgList) {
                    String msgBody = new String(msg.getBody(), Charset.forName(RemotingHelper.DEFAULT_CHARSET));
                    System.out.println("消费者A接收到消息:" +  "===== " + msgBody);
                    /*User user = JSON.parseObject(msgBody, User.class);
                    System.out.println("消费者A接收到消息:" +  "===== " + user.getId());*/
                }
                /*MessageExt messageExt = msgList.get(0);
                String body = new String(messageExt.getBody(), StandardCharsets.UTF_8);
                System.out.println("消费者A接收到消息: " + messageExt.toString() + "---消息内容为:" + body);*/
                // 标记该消息已经被成功消费
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });
        // 启动消费者实例
        mqPushConsumer.start();
        System.out.println("ConsumerA Started.");
    }
}

3.3 消费者B

package com.shucha.deveiface.biz.mq.consumer;
import com.shucha.deveiface.biz.constants.MqConstants;
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.common.protocol.heartbeat.MessageModel;
import java.util.List;
/**
 * @author tqf
 * @Description 消费者B
 * @Version 1.0
 * @since 2022-07-12 14:40
 */
public class ConsumerB {
    public static void main(String[] args) throws MQClientException {
        ConsumerB();
    }
    public static void ConsumerB() throws MQClientException {
        // 创建DefaultMQPushConsumer类并设定消费者名称
        DefaultMQPushConsumer mqPushConsumer = new DefaultMQPushConsumer(MqConstants.ConsumerGroup.CONSUMER_GROUP1);
        // 设置NameServer地址,如果是集群的话,使用分号;分隔开
        mqPushConsumer.setNamesrvAddr("127.0.0.1:9876");
        // 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费
        // 如果不是第一次启动,那么按照上次消费的位置继续消费
        mqPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
        // 设置消费模型,集群还是广播,默认为集群  CLUSTERING-集群   BROADCASTING-广播
        // mqPushConsumer.setMessageModel(MessageModel.CLUSTERING);
        mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);
        // 消费者最小线程量
        mqPushConsumer.setConsumeThreadMin(5);
        // 消费者最大线程量
        mqPushConsumer.setConsumeThreadMax(10);
        // 设置一次消费消息的条数,默认是1
        mqPushConsumer.setConsumeMessageBatchMaxSize(1);
        // 订阅一个或者多个Topic,以及Tag来过滤需要消费的消息,如果订阅该主题下的所有tag,则使用*
        mqPushConsumer.subscribe("TOPIC_TEST", "*");
        // 注册回调实现类来处理从broker拉取回来的消息
        mqPushConsumer.registerMessageListener(new MessageListenerConcurrently() {
            // 监听类实现MessageListenerConcurrently接口即可,重写consumeMessage方法接收数据
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgList, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
               /* MessageExt messageExt = msgList.get(0);
                String body = new String(messageExt.getBody(), StandardCharsets.UTF_8);
                System.out.println("消费者B接收到消息: " + messageExt.toString() + "---消息内容为:" + body);*/
                for (MessageExt msg : msgList) {
                    System.out.println("消费者B接收到消息:" +  "===== " + new String(msg.getBody()));
                    // String msgBody = new String(msg.getBody(), Charset.forName(RemotingHelper.DEFAULT_CHARSET));
                }
                // 标记该消息已经被成功消费
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });
        // 启动消费者实例
        mqPushConsumer.start();
        System.out.println("ConsumerB Started.");
    }
}

可以看到, 生产者发送了10条消息,ConsumerA与ConsumerB属于同一个消费者组,集群消费模式下每个消费者摊分消费所有消息。注意,两个消费者的ConsumerGroup组名需要一致,才算是同一个消费者组。

简单总结一下:

1、在Rocket集群消费模式下,(订阅)同一个主题(Topic)下的消息,对于不同的消费者组是一种“广播形式”,即每个消费者组的都会消费消息。

2、在Rocket集群消费模式下,(订阅)同一个主题(Topic)下的消息,对于相同的消费者组的消费者而言是一种集群模式,即同一个消费者组内的所有消费者均分消息并消费。

 4、广播消费

一条消息被多个 Consumer 消费,即使这些 Consumer 属于同一个 Consumer Group,消息也会被 Consumer Group 中的每个 Consumer 都消费一次,广播消费中的 Consumer Group 概念可以认为在消息划分方面无意 义。

使用方法:setMessageModel(MessageModel.BROADCASTING)

将前面的消费者A和消费者B的集群模式代码设置为如下

mqPushConsumer.setMessageModel(MessageModel.BROADCASTING);

重新启动生成者和2个消费者

4.1 生产者数据

4.2 消费者A

4.3 消费者B

可以看到生产者发送了10条消息,ConsumerA与ConsumerB属于同一个消费者组,广播模式下每个消费者都会全量消费所有消息 。

  • 集群消费:任何一条消息只需要被消费者集群中任意一个消费者处理。
  • 广播消费:每条消息被推送给消费者集群中的所有注册消费者,保证消息被每个消费者至少消费一次。 

到此这篇关于RocketMQ的两种消费模式详解的文章就介绍到这了,更多相关RocketMQ消费模式内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多支持脚本之家!

相关文章

  • Java基于享元模式实现五子棋游戏功能实例详解

    Java基于享元模式实现五子棋游戏功能实例详解

    这篇文章主要介绍了Java基于享元模式实现五子棋游戏功能,较为详细的分析了享元模式的概念、功能并结合实例形式详细分析了Java使用享元模式实现五子棋游戏的具体操作步骤与相关注意事项,需要的朋友可以参考下
    2018-05-05
  • Mybatis实现动态建表代码实例

    Mybatis实现动态建表代码实例

    这篇文章主要介绍了Mybatis实现动态建表代码实例,解释一下,就是指根据传入的表名,动态地创建数据库表,以供后面的业务场景使用,
    而使用 Mybatis 的动态 SQL,就能很好地为我们解决这个问题,需要的朋友可以参考下
    2023-10-10
  • Spring Boot如何优化内嵌的Tomcat示例详解

    Spring Boot如何优化内嵌的Tomcat示例详解

    spring boot默认web程序启用tomcat内嵌容器,监听8080端口,下面这篇文章主要给大家介绍了关于Spring Boot如何优化内嵌Tomcat的相关资料,文中通过示例代码介绍的非常详细,需要的朋友可以参考借鉴,下面来一起看看吧。
    2017-09-09
  • Java中的数组使用详解及练习

    Java中的数组使用详解及练习

    数组是Java程序中最常见的一种数据结构,它能够将相同类型的数据用一个标识符封装到一起,构成一个对象序列或基本数据类型,这篇文章主要给大家介绍了关于Java中数组使用详解及练习的相关资料,需要的朋友可以参考下
    2024-03-03
  • SpringBoot整合Quartz实现动态配置的代码示例

    SpringBoot整合Quartz实现动态配置的代码示例

    这篇文章将介绍如何把Quartz定时任务做成接口,实现以下功能的动态配置添加任务,修改任务,暂停任务,恢复任务,删除任务,任务列表,任务详情,文章通过代码示例介绍的非常详细,需要的朋友可以参考下
    2023-07-07
  • 如果你想写自己的Benchmark框架(推荐)

    如果你想写自己的Benchmark框架(推荐)

    这篇文章主要介绍了如果你想写自己的Benchmark框架,本文通过给大家分享八条军规,帮助大家理解,需要的朋友可以参考下
    2020-07-07
  • java freemarker实现动态生成excel文件

    java freemarker实现动态生成excel文件

    这篇文章主要为大家详细介绍了java如何通过freemarker实现动态生成excel文件,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下
    2023-12-12
  • Springboot配置返回日期格式化五种方法详解

    Springboot配置返回日期格式化五种方法详解

    本文主要介绍了Springboot配置返回日期格式化五种方法详解,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧
    2022-07-07
  • Java实现插入排序实例

    Java实现插入排序实例

    这篇文章主要介绍了Java实现插入排序,实例分析了Java的插入排序原理与实现技巧,非常具有实用价值,需要的朋友可以参考下
    2015-02-02
  • 详谈java线程与线程、进程与进程间通信

    详谈java线程与线程、进程与进程间通信

    下面小编就为大家带来一篇详谈java线程与线程、进程与进程间通信。小编觉得挺不错的,现在就分享给大家,也给大家做个参考。一起跟随小编过来看看吧
    2017-04-04

最新评论