赞
踩
任何MQ产品都可能存在各种异常,这些异常可能导致消息无法被发送到Broker,或者消息无法被消费者接收到,因此大部分MQ产品都会提供消息失败的重试机制。RocketMQ也不例外,在RocketMQ中消息重试分为生产者端重试和消费者端重试两种类型。
生产者端重试是指当生产者向Broker发送消息时,如果当前网络抖动等原因导致消息发送失败,此时可以通过手动设置发送失败重试次数的方式让消息重发一次。
- @RunWith(SpringRunner.class)
- @SpringBootTest
- public class RockermqproducerApplicationTests {
-
-
- @Value("${apache.rocketmq.producer.producerGroup}")
- private String producerGroup;
-
-
- /**
- * NameServer 地址
- */
- @Value("${apache.rocketmq.namesrvAddr}")
- private String namesrvAddr;
-
-
- @Test
- public void contextLoads() {
-
- //生产者的组名
- DefaultMQProducer producer=new DefaultMQProducer(producerGroup);
- //指定NameServer地址,多个地址以 ; 隔开
- producer.setNamesrvAddr(namesrvAddr);
- //消息发送失败重试次数
- producer.setRetryTimesWhenSendFailed(3);
- //异步发送失败重试次数
- producer.setRetryTimesWhenSendAsyncFailed(3);
- //消息没有发送成功,是否发送到另外一个Broker中
- producer.setRetryAnotherBrokerWhenNotStoreOK(true);
- try {
-
- /**
- * Producer对象在使用之前必须要调用start初始化,初始化一次即可
- * 注意:切记不可以在每次发送消息时,都调用start方法
- */
- producer.start();
- for (int i=0;i<=10000;i++)
- {
- Message msg=new Message("topic_example_java","TagA",("Hello Java Demo RocketMQ:"+i).getBytes(Charset.defaultCharset()));
- SendResult result=producer.send(msg);
- System.out.println("消息发送结果:"+result);
- }
-
- }catch (Exception e)
- {
- e.printStackTrace();
- }finally {
- producer.shutdown();
- }
- }
-
- }
可以看到通过org.apache.rocketmq.client.producer.DefaultMQProducer类的setRetryTimesWhenSendFailed和setRetryTimesWhenSendAsyncFailed方法设置了重试次数。这里实现重试逻辑的代码主要在DefaultMQProducerImpl类的sendDefaultImpl方法中。
消费者端的失败一般分为两种情况,一种是由于网络等原因导致消息没法从Broker发送到消费者端,这时在RocketMQ内部会不断尝试发送这条消息,直到发送成功为止(比如向集群中的一个Broker实例发送失败,就尝试发往另一个Broker实例);二是消费者端已经正常接收到消息了,但是在执行后续消息处理逻辑时发生了异常,最终反馈给MQ消费者处理失败,例如所接收到的消息数据可能不符合本身的业务要求,如当前卡号未激活不能执行业务等,这时就需要通过业务代码返回消息消费的不同状态来控制。
接下来以普通消费为例,看一下当消费者端出现业务消息消费异常之后时如何进行重试的。下面死在消费者端代码中注册消息监听器的consumeMessage方法最终返回的消息消费状态ConsumeConcurrentlyStatus的定义:
- public enum ConsumeConcurrentlyStatus {
- CONSUME_SUCCESS,
- RECONSUME_LATER;
-
- private ConsumeConcurrentlyStatus() {
- }
- }
CONSUME_SUCCESS表示消费成功,这是正常业务代码中返回的状态。RECONSUME_LATER表示当前消费失败,需要稍后进行重试。在RocketMQ中只有业务消费者侧返了CONSUME_SUCCESS才会认为消息消费时成功的,如果返回RECONSUME_LATER,RocketMQ则会认为消费失败,需要重新投递。为了保证消息至少被成功消费一次,RocketMQ会把认为消费失败的消息发回Broker,在接下来的某个时间点(默认是10秒,可修改)再次投递给消费者。如果一直重复消息都失败的话,当失败累积到一定次数后(默认16次)将消息投递到死信队列(Dead Letter Queue)中,此时需要监控死信队列进行人工干预。
- @Bean("consumer1")
- public DefaultMQPushConsumer consumer1()
- {
- //创建一个消息消费者,并设置一个消息消费者组
- DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("niwei_consumer_group");
- //指定 NameServer 地址
- consumer.setNamesrvAddr("localhost:9876");
- consumer.setInstanceName("PushConsumer1");
- //设置 Consumer 第一次启动时从队列头部开始消费还是队列尾部开始消费
- consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
- try {
- //订阅指定 Topic 下的所有消息
- consumer.subscribe("topic_example_java", "*");
- //注册消息监听器
- consumer.registerMessageListener(new MessageListenerConcurrently() {
- public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext context) {
- //默认 list 里只有一条消息,可以通过设置参数来批量接收消息
- if (list != null) {
- for (MessageExt ext : list) {
- String msgBody=new String(ext.getBody());
- if (ext.getReconsumeTimes()==3)
- {
- saveReconsumeStillMessage(ext);
- return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
- }else
- {
- try {
- doBusiness(msgBody);
- }catch (Exception e)
- {
- return ConsumeConcurrentlyStatus.RECONSUME_LATER;
- }
- }
-
- }
- }
- return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
- }
- });
-
- // 消费者对象在使用之前必须要调用 start 初始化
- consumer.start();
- System.out.println("消息消费者已启动");
- }catch (Exception e)
- {
- e.printStackTrace();
- }
- return consumer;
- }
示例中,判断当前消息是经过重试3次之后发出的,则不再继续执行业务代码,直接记录消息数据并返回消费成功状态。如果早业务的回调中没有处理好异常返回状态,而是在方法执行过程中抛出异常,那么RocketMQ认为消费也是失败的,会当作RECONSUME_LATER来处理。
当使用顺序消费的回调(实现了MessageListenerOrderly接口)时,由于顺序消费是只有前一条消息消费成功了才能继续,所以在其消息状态定义(ConsumeOrderlyStatus)中并没有RECONSUME_LATER状态,而是用SUSPEND_CURRENT_QUEUE_A_MOMENT来暂停当前队列的消费动作,直到消息经过不断重试成功为止。
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。