当前位置: 首页 > news >正文

RabbitMQ高级特性----生产者确认机制

题记:在Java微服务开发中,对于一个功能需要调用另一个服务下的功能才能实现的情况,我们通常会使用异步调用取代同步调用,进而实现增强业务的可拓展性和实现故障隔离以及流量削峰填谷的目的。而消息队列就是异步调用的解决方案之一。不过在使用消息队列实现异步调用的时候,可能会出现消息无法传递到位进而导致业务信息出现差异的情况,因此消息的传递的可靠性就显得尤为重要。

传递消息的流程:

要保障消息传递的可靠性,我们可以从消息队列的每一步分析,以RabbitMQ为例,其发送消息的流程大致如下图所示:

不难发现消息传递一共会经历三名角色,分别是消息发送者,MQ和消息接收者。因此我们要保证消息成功传递并被正确处理就需要保证这三者的可靠性。

发送者可靠性:

1.1生产者重连

首先第一种情况,就是生产者发送消息时,出现了网络故障,导致与MQ的连接中断。为了解决这个问题,SpringAMQP提供的消息发送时的重试机制。即:当RabbitTemplate与MQ连接超时后,多次重试。

我们可以在消息发送者的yml配置文件中加入如下信息:

spring: rabbitmq: connection-timeout: 1s # 设置MQ的连接超时时间 template: retry: enabled: true # 开启超时重试机制,默认时关闭的 initial-interval: 1000ms # 失败后的初始等待时间 multiplier: 2 # 失败后下次的等待时长倍数,下次等待时长 = initial-interval * multiplier max-attempts: 3 # 最大重试次数

我们可以在消息发送者工程下创建一个Test模拟向rabbitMQ发送消息(在此之前我们需要把rabbiitMQ关闭或是断掉网络)

@Test public void TestTimeOutPublish(){ rabbitTemplate.convertAndSend("simple.queue","Hello,world"); }

代码执行后会是这样的效果:

通过日志我们可以观察到一共重连了三次,与我们在yml文件中配置的max-attempt属性一致。需要注意的是,这里等待的时间还需要再加上connect-timeout这一判定连接超时的配置。

1.2生产者确认机制

在保证网络畅通并且RabbitMQ服务正确启动了的前提下,我们就可以成功将消息发送到MQ中了。

不过仍有少数情况下会出现消息进入mq后丢失的情况:

  • 无法找到指定的exchange,大部分情况下就是交换机名称出错
  • exchange无法正确路由到queue,可能是routeKey错误导致的

针对上述情况,RabbitMQ提供了生产者消息确认机制,包括Publisher ConfirmPublisher Return两种。在开启确认机制的情况下,当生产者发送消息给MQ后,MQ会根据消息处理的情况返回不同的回执

当我们在yml文件中增加了相应的配置后流程就变得如下图所示:

配置完成之后:

  • 对于所有成功投递如Mq中的消息都会返回Ack,表示投递成功。
  • 对于投递成功但路由失败的消息,会返回Publish Return并返回Ack
  • 除此之外的所有消息都返回NAck,表示消息投递失败

需要注意的是,对于临时消息而言是进入队列Publisher Comfirm就是返回Ack,如果是持久消息则需要写入磁盘后才是返回Ack。

实现方式如下:

publisher消息发送者yml配置文件中加入:

spring: rabbitmq: publisher-confirm-type: correlated # 开启publisher confirm机制,并设置confirm类型 publisher-returns: true # 开启publisher return机制

pubulisher-confirm-type一共有三种属性提供选择:

  1. none:默认选项,关闭
  2. correlated:异步回调返回回执
  3. simple:同步阻塞等待mq回执

并且我们还需要在配置类中配置confrim Return报文:

public class MqConfig { private final RabbitTemplate rabbitTemplate; @PostConstruct //在注入rabbitTemplate依赖后执行 public void init(){ rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() { @Override public void returnedMessage(ReturnedMessage returned) { log.error("触发return callback,"); log.debug("exchange: {}", returned.getExchange()); log.debug("routingKey: {}", returned.getRoutingKey()); log.debug("message: {}", returned.getMessage()); log.debug("replyCode: {}", returned.getReplyCode()); log.debug("replyText: {}", returned.getReplyText()); } }); } }

对于不同的业务我们有不同的处理方案,因此CallbackComfirm需要在每次发送方法前定义,并且由作为convertMessage方法的参数,如下所示:

@Test void testPublisherConfirm() { // 1.创建CorrelationData CorrelationData cd = new CorrelationData(); // 2.给Future添加ConfirmCallback cd.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() { @Override public void onFailure(Throwable ex) { // 2.1.Future发生异常时的处理逻辑,基本不会触发 log.error("send message fail", ex); } @Override public void onSuccess(CorrelationData.Confirm result) { // 2.2.Future接收到回执的处理逻辑,参数中的result就是回执内容 if(result.isAck()){ // result.isAck(),boolean类型,true代表ack回执,false 代表 nack回执 log.debug("发送消息成功,收到 ack!"); }else{ // result.getReason(),String类型,返回nack时的异常描述 log.error("发送消息失败,收到 nack, reason : {}", result.getReason()); } } }); // 3.发送消息 rabbitTemplate.convertAndSend("simple.direct", "simple", "hello,RabbiteMQ", cd); }

我们可以在onSuccess和onFailure中编写我们对这一消息发送成功以及失败情况的对应处理。

如果关于上述有其他更好的建议以及疑问,欢迎留言,我将尽快回复。

http://icebutterfly214.com/news/240394/

相关文章:

  • 大数据预测分析在餐饮行业的市场趋势预测
  • 【教程4>第10章>第20节】基于FPGA的图像sobel锐化算法开发——图像sobel锐化仿真测试以及MATLAB辅助验证
  • Java Web 教学资源库系统源码-SpringBoot2+Vue3+MyBatis-Plus+MySQL8.0【含文档】
  • MPC5634 Bootloader
  • WSL Ubuntu 安装 Docker 操作指南
  • 【教程4>第10章>第21节】基于FPGA的图像Laplace边缘提取算法开发——理论分析与matlab仿真
  • 前后端分离BB平台系统|SpringBoot+Vue+MyBatis+MySQL完整源码+部署教程
  • 【毕业设计】SpringBoot+Vue+MySQL 知识管理系统平台源码+数据库+论文+部署文档
  • SpringBoot+Vue 高校学科竞赛平台管理平台源码【适合毕设/课设/学习】Java+MySQL
  • Linux命令-ipcrm命令(删除Linux系统中的进程间通信(IPC)资源)
  • openmv与stm32通信入门必看:手把手教程(从零实现)
  • [内网流媒体] 能长期使用的内网工具具备哪些特征
  • 零基础学习Proteus元器件库大全与原理图绘制流程
  • 零基础学习JLink下载的完整操作流程
  • DeviceMetadataParsers.dll文件丢失找不到问题 免费下载方法分享
  • d3d10_1core.dll文件丢失找不到 彻底修复解决办法分享
  • ego1开发板大作业vivado实践指南:温度传感器数据采集系统
  • Day 09:【99天精通Python】字典与集合 - 键值对与去重利器
  • ModbusTCP协议详解实时性优化在STM32上的实践
  • 强化学习算法
  • Redis集群:原理与实战经验分享(面试必看!)
  • 计算机毕设 java 基于 Android 的医疗预约系统的设计与实现 移动医疗预约服务平台 医患对接信息化系统
  • 基于Java+SpringBoot+SSM乡村支教管理系统(源码+LW+调试文档+讲解等)/乡村教育支援系统/支教管理平台/乡村支教项目系统/农村支教管理系统/支教信息管理系统/乡村教师支援系统
  • 一文说清image2lcd图像转换核心要点
  • 深度解析:AI提示系统技术架构中的多轮对话管理设计
  • STM32下vTaskDelay实现任务延时的完整指南
  • 双主模式I2C在工业系统中的应用:完整示例
  • STM32CubeMX初学者指南:零基础快速理解开发流程
  • vivado2022.2安装全流程图文并茂的系统学习资料
  • ⚡_实时系统性能优化:从毫秒到微秒的突破[20260110165821]