RocketMQ重试机制及消息幂代码实例解析

这篇文章主要介绍了RocketMQ重试机制及消息幂代码实例解析,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友可以参考下

一.重试机制

  1.由于MQ经常处于复杂的分布式系统中,考虑网络波动,服务宕机,程序异常因素,很有可能出现消息发送或者消费失败的问题。因此,消息的重试就是所有MQ中间件必须考虑到的一个关键点。如果没有消息重试,就可能产生消息丢失的问题,可能对系统产生很大的影响。所以,秉承宁可多发消息,也不可丢失消息的原则,大部分MQ都对消息重试提供了很好的支持。

  2.RocketMQ为了使用者封装了消息重试的处理流程,无需开发人员手动处理。RocketMQ支持了生产端和消费端两类重试机制。

模拟异常

  Consumer端消息消费两种状态:

package com.alibaba.rocketmq.client.consumer.listener;

public enum ConsumeConcurrentlyStatus {
  CONSUME_SUCCESS,
  RECONSUME_LATER;

  private ConsumeConcurrentlyStatus() {
  }
}

  一个是成功(CONSUME_SUCCESS),一个是失败&重试(RECONSUME_LATER);

  Consumer为了保证消息消费成功,只有使用方明确表示消费成功,返回CONSUME_SUCCESS,RocketMQ才会认为消息消费成功。

  如果消息消费失败,只要返回ConsumeConcurrentlyStatus.RECONSUME_LATER,RocketMQ就会认为消息消费失败了,需要重新投递。

  1.出现异常

package com.wn.consumer;

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;

public class MQConsumer {
  public static void main(String[] args) throws MQClientException {
    //创建消费者
    DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("rmq-group");
    //设置NameServer地址
    consumer.setNamesrvAddr("192.168.138.187:9876;192.168.138.188:9876");
    //设置消费者实例名称
    consumer.setInstanceName("consumer");
    //订阅topic
    consumer.subscribe("wn02","TagA");
    //监听消息
    consumer.registerMessageListener(new MessageListenerConcurrently() {
      @Override
      public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
        //获取消息
        for (MessageExt msg:list){
          System.out.println(msg.getMsgId()+"---"+new String(msg.getBody()));
        }

        try {
          int i=1/0;
        }catch (Exception e){
          e.printStackTrace();
          //需要重试
          return ConsumeConcurrentlyStatus.RECONSUME_LATER;
        }//消息成功
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
      }
    });
    consumer.start();
    System.out.println("Consumer Started...");

  }
}

  2.网络延迟

package com.wn.consumer;

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;

public class MQConsumer {
  public static void main(String[] args) throws MQClientException {
    //创建消费者
    DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("rmq-group");
    //设置NameServer地址
    consumer.setNamesrvAddr("192.168.138.187:9876;192.168.138.188:9876");
    //设置消费者实例名称
    consumer.setInstanceName("consumer");
    //订阅topic
    consumer.subscribe("wn03","TagA");
    //监听消息
    consumer.registerMessageListener(new MessageListenerConcurrently() {
      @Override
      public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
        //获取消息
        for (MessageExt msg:list){
          System.out.println(msg.getMsgId()+"---"+new String(msg.getBody()));
        }

        //网络延迟
        try {
          Thread.sleep(600000);
        } catch (InterruptedException e) {
          e.printStackTrace();
        }
//消息成功
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
      }
    });
    consumer.start();
    System.out.println("Consumer Started...");

  }
}

二、消息幂等

1、在什么情况下会发生RocketMQ的消息重复消费

   ①、当系统的调用链路比较长的时候,比如系统A调用系统B,系统B再把消息发送到RocketMQ中,在系统A调用系统B的时候,如果系统B处理成功,但是迟迟没有将调用成功的结果返回给系统A的时候,系统A就会尝试重新发起请求给系统B,造成系统B重复处理,发起多条消息给RocketMQ造成重复消费;

  ②、在系统B发送给RocketMQ的时候,也有可能会发生和上面一样的问题,消息发送超时,节骨系统B重试,导致RocketMQ接收到了重读消息;

  ③、当RocketMQ成功接收到消息,并将消息交给消费者处理,如果消费者消费完成后还没来得及提交offset给RocketMQ,自己宕机或者重启了,那么RocketMQ没有接收到offset,就会认为消费失败了,会重发消息给消费者再次消费;

2、如何解决消息的重复消费

  通过幂等性来保证,只要保证重复消息不对结果产生影响,就完美地解决这个问题。

在生产者端保证幂等性,一下两种方式:

  ①、RocketMQ支持消息查询的功能,只要去RocketMQ查询一下是否已经发送过该条消息就可以了,不存在则发送,存在则不发送;

  ②、引入Redis,在发送消息到RocketMQ成功之后,向Redis中插入一条数据,如果发送重试,则先去Redis查询一个该条消息是否已经发送过了,存在的话就不重复发送消息了;

  方法一:RocketMQ消息查询的性能不是特别好,如果在高并发的场景下,每条消息在发送到RocketMQ时都去查询一下,可能会影响接口的性能;

  方法二:在一些极端的场景下,Redis也无法保证消息发送成功之后,就一定能写入Redis成功,比如写入消息成功而Redis此时宕机,那么再次查询Redis判断消息是否已经发送过,是无法得到正确结果的;

3、生产者

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();
  }
}

4、消费者

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!");
  }
}

以上就是本文的全部内容,希望对大家的学习有所帮助,也希望大家多多支持我们。

(0)

相关推荐

  • 使用Kotlin+RocketMQ实现延时消息的示例代码

    一. 延时消息 延时消息是指消息被发送以后,并不想让消费者立即拿到消息,而是等待指定时间后,消费者才拿到这个消息进行消费. 使用延时消息的典型场景,例如: 在电商系统中,用户下完订单30分钟内没支付,则订单可能会被取消. 在电商系统中,用户七天内没有评价商品,则默认好评. 这些场景对应的解决方案,包括: 轮询遍历数据库记录 JDK 的 DelayQueue ScheduledExecutorService 基于 Quartz 的定时任务 基于 Redis 的 zset 实现延时队列. 除此之外,

  • Spring Boot优雅使用RocketMQ的方法实例

    前言 MQ,是一种跨进程的通信机制,用于上下游传递消息.在传统的互联网架构中通常使用MQ来对上下游来做解耦合. 举例:当A系统对B系统进行消息通讯,如A系统发布一条系统公告,B系统可以订阅该频道进行系统公告同步,整个过程中A系统并不关系B系统会不会同步,由订阅该频道的系统自行处理. 什么是RocketMQ?# 官方说明: 随着使用越来越多的队列和虚拟主题,ActiveMQ IO模块遇到了瓶颈.我们尽力通过节流,断路器或降级来解决此问题,但效果不佳.因此,我们那时开始关注流行的消息传递解决方案Ka

  • springBoot整合RocketMQ及坑的示例代码

    版本: JDK:1.8 springBoot:1.5.10 rocketMQ:4.2.0 pom 配置: <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>1.5.10.RELEASE</version> </parent> <d

  • Java RocketMQ 路由注册与删除的实现

    简介 RocketMQ路由注册与删除是通过Broker与NameServer的心跳功能实现的.Broker启动时向集群中所有的NameServer发送心跳语句,每隔30s向集群中所有NameServer发送心跳包,NameServer收到Broker心跳包时会更新brokerLiveTable中的lastUpdateTimestamp,然后NameServer每隔10s扫描brokerLiveTable,如果连续120s没有收到心跳包,NameServer将移除该Broker的路由信息. 路由信

  • rocketmq消费负载均衡--push消费详解

    前言 本文介绍了DefaultMQPushConsumerImpl消费者,客户端负载均衡相关知识点.本文从DefaultMQPushConsumerImpl启动过程到实现负载均衡,从源代码一步一步分析,共分为6个部分进行介绍,其中第6个部分 rebalanceByTopic 为负载均衡的核心逻辑模块,具体过程运用了图文进行阐述. 介绍之前首先抛出几个问题: 1. 要做负载均衡,首先要解决的一个问题是什么? 2. 负载均衡是Client端处理还是Broker端处理? 个人理解: 1. 要做负载均衡

  • 浅谈Springboot整合RocketMQ使用心得

    一.阿里云官网---帮助文档 https://help.aliyun.com/document_detail/29536.html?spm=5176.doc29535.6.555.WWTIUh 按照官网步骤,创建Topic.申请发布(生产者).申请订阅(消费者) 二.代码 1.配置: public class MqConfig { /** * 启动测试之前请替换如下 XXX 为您的配置 */ public static final String PUBLIC_TOPIC = "test"

  • linux安装RocketMQ实例步骤

    1.安装JDK 1.1 检查当前虚拟机环境有没有JDK   rpm -qa|grep java 1.2 卸载  rpm -e --nodeps xxxxxx(自己的openjdk) 1.3 安装JDK 在/usr/local新建一个java文件夹,然后将tar包上传到文件夹下 切换到/usr/local/java   使用tar  -zxvf xxx解压 配置/etc/profile文件,加入JDK环境变量 export JAVA_HOME=/usr/local/java/jdk1.8.0_12

  • Window搭建部署RocketMQ步骤详解

    序 以前简单用过ActiveMQ但是公司项目上使用的是RocketMQ,所以准备多花点时间在这上面,搞懂项目的配置使用. 看了很多资料,先说说我自己对RocketMQ的简单理解.不管是我们写的消费者还是生产者都属于客户端,而我们需要安装RocketMQ,这是属于服务端.和ActivieMQ.zookeeper类似,消费者.生成者.服务端(NameServer)之间是采取观察者模式实现. 在操作系统上安装RocketMQ,启动服务端NameServer.启动Broker,书写Consumer代码,

  • Docker中RocketMQ的安装与使用详解

    搜索RocketMQ的镜像,可以通过docker的hub.docker.com上进行搜索,也可以在Linux下通过docker的search命令进行搜索,不过最近防火墙升级后,导致国外的网站打开都很慢,通过命令搜索反而会更加方便,操作Docker命令一定要是root用户或者具有root权限的用户.查询操作如下: docker search rocketmq 可以得到如下的结果: 镜像倒是蛮多的,不过看来看去没有一个是官方发布的,我就随便选一个吧,如foxiswho/rocketmq,以下是一个查

  • java RocketMQ快速入门基础知识

    如何使用 1.引入 rocketmq-client <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>4.1.0-incubating</version> </dependency> 2.编写Producer DefaultMQProducer produce

随机推荐