RocketMQ入门案例

  学习RocketMQ,先写一个Demo演示一下看看效果。

一、服务端部署

  因为只是简单的为了演示效果,服务端仅部署单Master模式 —— 一个Name Server节点,一个Broker节点。主要有以下步骤。

  1. 下载RocketMQ源码、编译(也可以网上下载编译好的文件),这里使用最新的4.4.0版本,下载好之后放在Linux上通过一下命令解压缩、编译。

    1unzip rocketmq-all-4.4.0-source-release.zip 2cd rocketmq-all-4.4.0/ 3mvn -Prelease-all -DskipTests clean install –U
  2. 编译之后到distribution/target/apache-rocketmq目录,后续所有操作都是在该路径下。

    cd distribution/target/apache-rocketmq
    
  3. 启动Name Server,查看日志确认启动成功。

    1nohup sh bin/mqnamesrv & 2tail -f ~/logs/rocketmqlogs/namesrv.log
  4. 启动Broker,查看日志确认启动成功。

    1nohup sh bin/mqbroker -n localhost:9876 & 2tail -f ~/logs/rocketmqlogs/broker.log

  Name Server和Broker都成功启动,服务器就部署完成了。更详细的参考官方文档手册,里面还包含在服务器上运行Producer、Customer示例,这里主要是在项目中使用。

  官网手册戳这里:Quick Start

二、客户端搭建:Spring Boot项目中使用

  客户端分为消息生产者和消息消费者,这里通过日志打印输出查看效果,为了看起来更清晰,我新建了两个模块分别作为消息生产者和消息消费者。

  1. 添加依赖,在两个模块的pom文件中添加以下配置。

    1<dependency> 2 <groupId>org.apache.rocketmq</groupId> 3 <artifactId>rocketmq-client</artifactId> 4 <version>4.4.0</version> 5</dependency>
  2. 配置生产者模块。

    • application.yml文件中增加用来初始化producer的相关配置,这里只配了一部分,更详细的配置参数可以查看官方文档。

      1# RocketMQ生产者 2rocketmq: 3 producer: 4 # Producer组名,多个Producer如果属于一个应用,发送同样的消息,则应该将它们归为同一组。默认DEFAULT_PRODUCER 5 producerGroup: ${spring.application.name} 6 # namesrv地址 7 namesrvAddr: 192.168.101.213:9876 8 # 客户端限制的消息大小,超过报错,同时服务端也会限制,需要跟服务端配合使用。默认4MB 9 maxMessageSize: 4096 10 # 发送消息超时时间,单位毫秒。默认10000 11 sendMsgTimeout: 5000 12 # 如果消息发送失败,最大重试次数,该参数只对同步发送模式起作用。默认2 13 retryTimesWhenSendFailed: 2 14 # 消息Body超过多大开始压缩(Consumer收到消息会自动解压缩),单位字节。默认4096 15 compressMsgBodyOverHowmuch: 4096 16 # 在发送消息时,自动创建服务器不存在的topic,需要指定Key,该Key可用于配置发送消息所在topic的默认路由。 17 createTopicKey: XIAO_LIU
    • 新增producer配置类,系统启动时读取yml文件的配置信息初始化producer。集群模式下,如果在同一个jvm中,要往多个的MQ集群发送消息,则需要创建多个的producer并设置不同的instanceName,默认不需要设置该参数。

      1@Configuration 2public class ProducerConfiguration { 3 private static final Logger LOGGER = LoggerFactory.getLogger(ProducerConfiguration.class); 4 5 /** 6 * Producer组名,多个Producer如果属于一个应用,发送同样的消息,则应该将它们归为同一组。默认DEFAULT_PRODUCER 7 */ 8 @Value("${rocketmq.producer.producerGroup}") 9 private String producerGroup; 10 /** 11 * namesrv地址 12 */ 13 @Value("${rocketmq.producer.namesrvAddr}") 14 private String namesrvAddr; 15 /** 16 * 客户端限制的消息大小,超过报错,同时服务端也会限制,需要跟服务端配合使用。默认4MB 17 */ 18 @Value("${rocketmq.producer.maxMessageSize}") 19 private Integer maxMessageSize; 20 /** 21 * 发送消息超时时间,单位毫秒。默认10000 22 */ 23 @Value("${rocketmq.producer.sendMsgTimeout}") 24 private Integer sendMsgTimeout; 25 /** 26 * 如果消息发送失败,最大重试次数,该参数只对同步发送模式起作用。默认2 27 */ 28 @Value("${rocketmq.producer.retryTimesWhenSendFailed}") 29 private Integer retryTimesWhenSendFailed; 30 /** 31 * 消息Body超过多大开始压缩(Consumer收到消息会自动解压缩),单位字节。默认4096 32 */ 33 @Value("${rocketmq.producer.compressMsgBodyOverHowmuch}") 34 private Integer compressMsgBodyOverHowmuch; 35 /** 36 * 在发送消息时,自动创建服务器不存在的topic,需要指定Key,该Key可用于配置发送消息所在topic的默认路由。 37 */ 38 @Value("${rocketmq.producer.createTopicKey}") 39 private String createTopicKey; 40 41 @Bean 42 public DefaultMQProducer getRocketMQProducer() { 43 44 DefaultMQProducer producer = new DefaultMQProducer(this.producerGroup); 45 producer.setNamesrvAddr(this.namesrvAddr); 46 producer.setCreateTopicKey(this.createTopicKey); 47 48 if (this.maxMessageSize != null) { 49 producer.setMaxMessageSize(this.maxMessageSize); 50 } 51 if (this.sendMsgTimeout != null) { 52 producer.setSendMsgTimeout(this.sendMsgTimeout); 53 } 54 if (this.retryTimesWhenSendFailed != null) { 55 producer.setRetryTimesWhenSendFailed(this.retryTimesWhenSendFailed); 56 } 57 if (this.compressMsgBodyOverHowmuch != null) { 58 producer.setCompressMsgBodyOverHowmuch(this.compressMsgBodyOverHowmuch); 59 } 60 if (Strings.isNotBlank(this.createTopicKey)) { 61 producer.setCreateTopicKey(this.createTopicKey); 62 } 63 64 try { 65 producer.start(); 66 67 LOGGER.info("Producer Started : producerGroup:[{}], namesrvAddr:[{}]" 68 , this.producerGroup, this.namesrvAddr); 69 } catch (MQClientException e) { 70 LOGGER.error("Producer Start Failed : {}", e.getMessage(), e); 71 } 72 return producer; 73 } 74 75}
    • 使用producer实例向MQ发送消息。

      1@RunWith(SpringRunner.class) 2@SpringBootTest 3public class ProducerServiceApplicationTests { 4 private static final Logger LOGGER = LoggerFactory.getLogger(ProducerServiceApplicationTests.class); 5 @Autowired 6 private DefaultMQProducer defaultMQProducer; 7 8 @Test 9 public void send() throws MQClientException, RemotingException, MQBrokerException, InterruptedException, UnsupportedEncodingException { 10 for (int i = 0; i < 100; i++) { 11 User user = new User(); 12 user.setUsername("用户" + i); 13 user.setPassword("密码" + i); 14 user.setSex(i % 2); 15 user.setBirthday(new Date()); 16 Message message = new Message("user-topic", "user-tag", JSON.toJSONString(user).getBytes(RemotingHelper.DEFAULT_CHARSET)); 17 SendResult sendResult = defaultMQProducer.send(message); 18 LOGGER.info(sendResult.toString()); 19 } 20 } 21}
  3. 配置消费者模块。

    • application.yml文件中增加用来初始化consumer的相关配置,同样参数这里只配了一部分,更详细的配置参数可以查看官方文档。

      1# RocketMQ消费者 2rocketmq: 3 consumer: 4 # Consumer组名,多个Consumer如果属于一个应用,订阅同样的消息,且消费逻辑一致,则应该将它们归为同一组。默认DEFAULT_CONSUMER 5 consumerGroup: ${spring.application.name} 6 # namesrv地址 7 namesrvAddr: 192.168.101.213:9876 8 # 消费线程池最大线程数。默认10 9 consumeThreadMin: 10 10 # 消费线程池最大线程数。默认20 11 consumeThreadMax: 20 12 # 批量消费,一次消费多少条消息。默认1 13 consumeMessageBatchMaxSize: 1 14 # 批量拉消息,一次最多拉多少条。默认32 15 pullBatchSize: 32 16 # 订阅的主题 17 topics: user-topic
    • 新增consumer配置。

      1@Configuration 2public class ConsumerConfiguration { 3 private static final Logger LOGGER = LoggerFactory.getLogger(ConsumerConfiguration.class); 4 5 @Value("${rocketmq.consumer.consumerGroup}") 6 private String consumerGroup; 7 @Value("${rocketmq.consumer.namesrvAddr}") 8 private String namesrvAddr; 9 @Value("${rocketmq.consumer.consumeThreadMin}") 10 private int consumeThreadMin; 11 @Value("${rocketmq.consumer.consumeThreadMax}") 12 private int consumeThreadMax; 13 @Value("${rocketmq.consumer.consumeMessageBatchMaxSize}") 14 private int consumeMessageBatchMaxSize; 15 @Value("${rocketmq.consumer.pullBatchSize}") 16 private int pullBatchSize; 17 @Value("${rocketmq.consumer.topics}") 18 private String topics; 19 20 private final ConsumeMsgListener consumeMsgListener; 21 22 @Autowired 23 public ConsumerConfiguration(final ConsumeMsgListener consumeMsgListener) { 24 this.consumeMsgListener = consumeMsgListener; 25 } 26 27 @Bean 28 public DefaultMQPushConsumer getRocketMQConsumer() { 29 DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup); 30 consumer.setNamesrvAddr(namesrvAddr); 31 consumer.setConsumeThreadMin(consumeThreadMin); 32 consumer.setConsumeThreadMax(consumeThreadMax); 33 consumer.setConsumeMessageBatchMaxSize(consumeMessageBatchMaxSize); 34 consumer.setPullBatchSize(pullBatchSize); 35 consumer.registerMessageListener(consumeMsgListener); 36 37 try { 38 /** 39 * 设置消费者订阅的主题和tag。subExpression参数为*表示订阅该主题下所有tag, 40 * 如果需要订阅该主题下的指定tag,subExpression设置为对应tag名称,多个tag以||分割,例如"tag1 || tag2 || tag3" 41 */ 42 consumer.subscribe(topics, "*"); 43 consumer.start(); 44 45 LOGGER.info("Consumer Started : consumerGroup:{}, topics:{}, namesrvAddr:{}", consumerGroup, topics, namesrvAddr); 46 } catch (Exception e) { 47 LOGGER.error("Consumer Start Failed : consumerGroup:{}, topics:{}, namesrvAddr:{}", consumerGroup, topics, namesrvAddr, e); 48 e.printStackTrace(); 49 } 50 return consumer; 51 } 52}
    • 新增消息监听器,监听到新消息后,执行对应的业务逻辑。

      1@Component 2public class ConsumeMsgListener implements MessageListenerConcurrently { 3 private static final Logger LOGGER = LoggerFactory.getLogger(ConsumeMsgListener.class); 4 5 @Override 6 public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) { 7 if (CollectionUtils.isEmpty(msgs)) { 8 LOGGER.info("Msgs is Empty."); 9 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; 10 } 11 for (MessageExt msg : msgs) { 12 try { 13 if ("user-topic".equals(msg.getTopic())) { 14 LOGGER.info("{} Receive New Messages: {}", Thread.currentThread().getName(), new String(msg.getBody())); 15 // do something 16 } 17 } catch (Exception e) { 18 if (msg.getReconsumeTimes() == 3) { 19 // 超过3次不再重试 20 LOGGER.error("Msg Consume Failed."); 21 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; 22 } else { 23 // 重试 24 return ConsumeConcurrentlyStatus.RECONSUME_LATER; 25 } 26 } 27 } 28 29 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; 30 } 31}

三、效果

  1. 运行生产者测试代码。系统启动时初始化Producer,然后执行测试代码,往MQ中发送消息。效果如下:

  2. 启动消费者服务。系统启动时先初始化Customer。此时1.已经往MQ中发送了一些消息,监听器监听到MQ中有消息,随即马上消费消息。

四、总结

  Demo很简单,但是里面还有很多东西需要慢慢研究。

  代码可以戳这里:spring-cloud-learn

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang