因为MQ常常处于复杂的分布式系统中,考虑网络波动,服务宕机,程序异常因素,颇有可能出现消息发送或者消费失败的问题。所以,消息的重试就是全部MQ中间件必须考虑到的一个关键点。若是没有消息重试,就可能产生消息丢失的问题,可能对系统产生很大的影响。因此,秉承宁肯多发消息,也不可丢失消息的原则,大部分MQ都对消息重试提供了很好的支持。java
MQ 消费者的消费逻辑失败时,能够经过设置返回状态达到消息重试的结果。redis
MQ 消息重试只针对集群消费方式生效;广播方式不提供失败重试特性,即消费失败后,失败消息再也不重试,继续消费新的消息。网络
消费者:并发
package com.zn.retry; import com.alibaba.rocketmq.client.consumer.DefaultMQPushConsumer; import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import com.alibaba.rocketmq.client.consumer.listener.MessageListenerConcurrently; import com.alibaba.rocketmq.client.exception.MQClientException; import com.alibaba.rocketmq.common.message.MessageExt; import java.util.List; /** * RocketMQ重试机制消费者 */ public class RetryConsumer { public static void main(String[] args) throws MQClientException { //建立消费者 DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("rmq-group"); //设置NameServer地址 consumer.setNamesrvAddr("192.168.33.135:9876;192.168.33.136:9876"); //设置实例名称 consumer.setInstanceName("consumer"); //订阅topic consumer.subscribe("itmayiedu-topic","TagA"); //监听消息 consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) { //获取消息 for (MessageExt messageExt:list){ System.out.println(messageExt.getMsgId()+"---"+new String(messageExt.getBody())); } try { //模拟错误 int i=5/0; }catch (Exception e){ e.printStackTrace(); //须要重试 return ConsumeConcurrentlyStatus.RECONSUME_LATER; } //不须要重试 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); //启动消费者 consumer.start(); System.out.println("Consumer Started!"); } }
控制台效果:分布式
package com.zn.retry; import com.alibaba.rocketmq.client.consumer.DefaultMQPushConsumer; import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import com.alibaba.rocketmq.client.consumer.listener.MessageListenerConcurrently; import com.alibaba.rocketmq.client.exception.MQClientException; import com.alibaba.rocketmq.common.message.MessageExt; import java.util.List; /** * RocketMQ重试机制消费者 */ public class RetryConsumer { public static void main(String[] args) throws MQClientException { //建立消费者 DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("rmq-group"); //设置NameServer地址 consumer.setNamesrvAddr("192.168.33.135:9876;192.168.33.136:9876"); //设置实例名称 consumer.setInstanceName("consumer"); //订阅topic consumer.subscribe("itmayiedu-topic","TagA"); //监听消息 consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) { //获取消息 for (MessageExt messageExt:list){ System.out.println(messageExt.getMsgId()+"---"+new String(messageExt.getBody())); } try { //网络延迟 Thread.sleep(600000); } catch (InterruptedException e) { e.printStackTrace(); } //消费成功 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); //启动消费者 consumer.start(); System.out.println("Consumer Started!"); } }
①、当系统的调用链路比较长的时候,好比系统A调用系统B,系统B再把消息发送到RocketMQ中,在系统A调用系统B的时候,若是系统B处理成功,可是迟迟没有将调用成功的结果返回给系统A的时候,系统A就会尝试从新发起请求给系统B,形成系统B重复处理,发起多条消息给RocketMQ形成重复消费;ide
②、在系统B发送给RocketMQ的时候,也有可能会发生和上面同样的问题,消息发送超时,节骨系统B重试,致使RocketMQ接收到了重读消息;高并发
③、当RocketMQ成功接收到消息,并将消息交给消费者处理,若是消费者消费完成后还没来得及提交offset给RocketMQ,本身宕机或者重启了,那么RocketMQ没有接收到offset,就会认为消费失败了,会重发消息给消费者再次消费;性能
经过幂等性来保证,只要保证重复消息不对结果产生影响,就完美地解决这个问题。spa
在生产者端保证幂等性,一下两种方式:code
①、RocketMQ支持消息查询的功能,只要去RocketMQ查询一下是否已经发送过该条消息就能够了,不存在则发送,存在则不发送;
②、引入Redis,在发送消息到RocketMQ成功以后,向Redis中插入一条数据,若是发送重试,则先去Redis查询一个该条消息是否已经发送过了,存在的话就不重复发送消息了;
方法一:RocketMQ消息查询的性能不是特别好,若是在高并发的场景下,每条消息在发送到RocketMQ时都去查询一下,可能会影响接口的性能;
方法二:在一些极端的场景下,Redis也没法保证消息发送成功以后,就必定能写入Redis成功,好比写入消息成功而Redis此时宕机,那么再次查询Redis判断消息是否已经发送过,是没法获得正确结果的;
package com.zn.idempotent; import com.alibaba.rocketmq.client.exception.MQBrokerException; import com.alibaba.rocketmq.client.exception.MQClientException; import com.alibaba.rocketmq.client.producer.DefaultMQProducer; import com.alibaba.rocketmq.client.producer.SendResult; import com.alibaba.rocketmq.common.message.Message; import com.alibaba.rocketmq.remoting.exception.RemotingException; /** * 消息幂等生产者 */ public class IdempotentProvider { public static void main(String[] args) throws MQClientException, InterruptedException, RemotingException, MQBrokerException { //建立一个生产者 DefaultMQProducer producer=new DefaultMQProducer("rmq-group"); //设置NameServer地址 producer.setNamesrvAddr("192.168.33.135:9876;192.168.33.136:9876"); //设置生产者实例名称 producer.setInstanceName("producer"); //启动生产者 producer.start(); //发送消息 for (int i=1;i<=1;i++){ //模拟网络延迟,每秒发送一次MQ Thread.sleep(1000); //建立消息,topic主题名称 tags临时值表明小分类, body表明消息体 Message message=new Message("itmayiedu-topic03","TagA",("itmayiedu-"+i).getBytes()); //消息的惟一标识 message.setKeys("订单消息:"+i); //发送消息 SendResult sendResult=producer.send(message); System.out.println("信息幂等问题来了:"+sendResult.toString()); } producer.shutdown(); } }
package com.zn.idempotent; import com.alibaba.rocketmq.client.consumer.DefaultMQPushConsumer; import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import com.alibaba.rocketmq.client.consumer.listener.MessageListenerConcurrently; import com.alibaba.rocketmq.client.exception.MQClientException; import com.alibaba.rocketmq.common.message.MessageExt; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.logging.LogManager; import java.util.logging.Logger; /** * 消息幂等消费者 */ public class IdempotentConsumer { static private Map<String, Object> logMap = new HashMap<>(); public static void main(String[] args) throws MQClientException { //建立消费者 DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("rmq-group"); //设置NameServer地址 consumer.setNamesrvAddr("192.168.33.135:9876;192.168.33.136:9876"); //设置实例名称 consumer.setInstanceName("consumer"); //订阅topic consumer.subscribe("itmayiedu-topic03","TagA"); //监听消息 consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) { String key=null; String msgId=null; for (MessageExt messageExt:list){ key=messageExt.getKeys(); //判读redis中有没有当前消息key if (logMap.containsKey(key)) { // 无需继续重试。 System.out.println("key:"+key+",已经消费,无需重试..."); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } //RocketMQ因为是集群环境,因此产生的消息ID可能会重复 msgId = messageExt.getMsgId(); System.out.println("key:" + key + ",msgid:" + msgId + "---" + new String(messageExt.getBody())); //将当前key保存在redis中 logMap.put(messageExt.getKeys(),messageExt); } try { int i=5/0; }catch (Exception e){ e.printStackTrace(); //人工补偿 return ConsumeConcurrentlyStatus.RECONSUME_LATER; } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); //启动消费者 consumer.start(); System.out.println("Consumer Started!"); } }