十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

RabbitMQ 基础篇:从零开始掌握消息队列核心概念

RabbitMQ 基础篇:从零开始掌握消息队列核心概念 目录一、课程背景二、初识 MQ同步调用 vs 异步调用2.1 同步调用2.2 异步调用2.3 小结三、MQ 技术选型RabbitMQ / ActiveMQ / RocketMQ / Kafka四、RabbitMQ 基本介绍4.1 整体架构与核心概念4.2 官方工作流概览五、SpringAMQP 快速入门5.1 引入依赖5.2 配置 RabbitMQ 服务端信息5.3 发送消息5.4 接收消息六、Work Queues工作队列6.1 任务模型6.2 消费者消息推送限制6.3 小结七、交换机类型7.1 Fanout 交换机广播7.2 Direct 交换机定向路由7.3 Topic 交换机话题八、声明队列和交换机8.1 基于 Bean 的方式8.2 基于 RabbitListener 注解的方式九、消息转换器十、写在最后一、课程背景我们以一个常见的登录场景为例用户登录 → 用户微服务查询并校验用户信息 → 风控微服务记录登录信息并判断登录风险 → 短信微服务发送风险短信 → 通知用户。如果登录事件希望做风控和短信提醒传统方式就是登录完成后依次调用风控、短信两个服务。但在真实业务里风控和短信通知跟登录主链路关系并不紧密完全可以异步处理——这就轮到 MQ消息队列登场了。整个 RabbitMQ 课程分为基础篇和高级篇基础篇高级篇同步和异步发送者重连MQ 技术选型发送者确认数据隔离MQ 持久化SpringAMQPLazyQueueWork 模式消费者确认MQ 消息转换器失败重试发布订阅模式业务幂等消息堆积问题处理延迟消息本文聚焦于基础篇的全部内容。二、初识 MQ同步调用 vs 异步调用2.1 同步调用以黑马商城的余额支付为例。一次支付要经过 5 个步骤扣减余额用户服务更新支付状态支付服务更新订单状态交易服务短信通知用户通知服务增加用户积分积分服务每一步 50ms串行下来一次请求需要300ms。这种链式同步调用会出现三个问题拓展性差每加一个下游业务就要改支付服务的代码性能下降链路越长整体 RT 越高级联失败任何一环出问题都会让整条链路失败。小结同步调用的优势是时效性强调用方等到结果后才返回缺点是拓展性差、性能下降、级联失败。2.2 异步调用异步调用本质上是基于消息通知的方式一般包含三个角色消息发送者投递消息的人就是原来的调用方消息代理管理、暂存、转发消息可以理解成消息服务器消息接收者接收并处理消息的人就是原来的服务提供方。用异步改造余额支付支付服务只同步做完扣减余额 更新支付状态这两件强关联的事剩下的订单状态、短信、积分通过 Broker 异步分发。改造后访问时长从 300ms 缩短到 100ms并且具备四个明显优势解除耦合拓展性强无需等待性能好故障隔离下游服务故障不影响上游业务缓存消息、流量削峰填谷把脉冲式的 QPS 削平成平缓曲线小结异步调用的优势是耦合度低、拓展性强、性能好、故障隔离、流量削峰缺点是不能立即得到调用结果时效性差不确定下游业务是否执行成功业务安全依赖于 Broker 的可靠性。2.3 小结维度同步调用异步调用时效性强等结果弱后置处理拓展性差强性能随链路下降稳定耦合度高低故障影响级联失败隔离流量处理难以削峰削峰填谷适用场景强一致性、强依赖最终一致性、可异步化的非关键路径⚠️什么时候用异步主链路完成后才需要的非关键业务风控、通知、积分、日志、审计…都适合异步化。但涉及金钱、需要强一致性的核心链路依然建议同步 本地事务。三、MQ 技术选型RabbitMQ / ActiveMQ / RocketMQ / KafkaMQMessageQueue中文消息队列本质上就是异步调用中的Broker。业界常见的有四款产品维度RabbitMQActiveMQRocketMQKafka公司/社区RabbitApache阿里Apache开发语言ErlangJavaJavaScala Java协议支持AMQP, XMPP, SMTP, STOMPOpenWire, STOMP, REST, XMPP, AMQP自定义协议自定义协议可用性高一般高高单机呑吐量一般差高非常高消息延迟微秒级毫秒级毫秒级毫秒以内消息可靠性高一般高一般怎么选强业务可靠、对延迟敏感RabbitMQ微秒级延迟 高可靠电商订单场景首选日志、大数据、流计算、IoT 等超高吞吐场景Kafka阿里系、有顺序消息、事务消息要求RocketMQ老系统维护、AMQP 协议强需求ActiveMQ。本课程选RabbitMQ作为讲解对象覆盖绝大多数业务场景。四、RabbitMQ 基本介绍4.1 整体架构与核心概念RabbitMQ 的整体架构里有五个核心概念先记牢virtual-host虚拟主机起到数据隔离的作用类似 MySQL 的 databasepublisher消息发送者consumer消息的消费者queue队列用于存储消息exchange交换机负责路由消息。publisher 把消息发给 exchangeexchange 按规则路由到一个或多个 queue最后由 consumer 消费。每个 VirtualHost 是一个独立的小型 RabbitMQ 实例。4.2 官方工作流概览Publisher → Exchange按 Binding 规则→ Queue → Consumer这是 RabbitMQ 的标准消息流。理解清楚exchange 只负责转发、不存储消息这一点对后面的代码实现非常关键。五、SpringAMQP 快速入门官方地址Spring AMQPSpring AMQP 是基于AMQP 协议定义的一套 API 规范包含两部分spring-amqp基础抽象spring-rabbit底层默认实现也就是我们用的 RabbitMQ。下面用 4 步把一个最简单的生产者发送字符串、消费者接收字符串跑起来。5.1 引入依赖在父工程中引入 spring-amqp 依赖这样 publisher 和 consumer 服务都可以使用dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency5.2 配置 RabbitMQ 服务端信息在每个微服务的application.yml中配置这样服务才能连到 RabbitMQspring: rabbitmq: host: 192.168.150.101 # 主机名 port: 5672 # 端口 virtual-host: /hmall # 虚拟主机 username: hmall # 用户名 password: 123 # 密码5.3 发送消息SpringAMQP 提供了RabbitTemplate工具类方便我们发送消息Autowired private RabbitTemplate rabbitTemplate; Test public void testSimpleQueue() { // 队列名称 String queueName simple.queue; // 消息 String message hello, spring amqp!; // 发送消息 rabbitTemplate.convertAndSend(queueName, message); }5.4 接收消息SpringAMQP 提供声明式的消息监听我们只需通过RabbitListener注解在方法上声明要监听的队列名称将来 SpringAMQP 就会自动把消息传递给当前方法Slf4j Component public class SpringRabbitListener { RabbitListener(queues simple.queue) public void listenSimpleQueueMessage(String msg) throws InterruptedException { log.info(spring 消费者接收到消息【 msg 】); if (true) { throw new MessageConversionException(故意的); } log.info(消息处理完成); } } 这一节抛出的MessageConversionException是为高级篇里讲失败重试 / 死信队列故意埋的伏笔基础篇先不用深究。六、Work Queues工作队列6.1 任务模型Work Queues又叫任务模型简单说就是让多个消费者绑定到一个队列共同消费队列中的消息。结构如下publisher → queue → consumer1 ↘ consumer2它的特点是同一消息只会被一个消费者处理通过多消费者能显著加快消息处理速度。6.2 消费者消息推送限制默认情况下RabbitMQ 会把消息依次轮询投递给绑定在队列上的每一个消费者。但它并不会等消费者处理完再投递下一条而是按平均分配消息数量的方式分发可能造成消息堆积。解决办法是修改application.yml设置prefetch值为 1确保同一时刻最多投递给消费者 1 条消息spring: rabbitmq: listener: simple: prefetch: 1 # 每次只能获取一条消息处理完成才能获取下一个消息这样无论是消费者 1 还是消费者 2处理完一条再处理下一条就能做到能者多劳——机器性能好的自然消费得快。6.3 小结Work 模型的使用要点多个消费者绑定到一个队列可以加快消息处理速度同一条消息只会被一个消费者处理通过prefetch控制消费者预取的消息数量处理完一条再处理下一条实现能者多劳。七、交换机类型真正生产环境里消息都是经过exchange发送的而不是直接发送到队列。交换机共有三种类型Fanout广播Direct定向Topic话题。7.1 Fanout 交换机广播Fanout Exchange 会将收到的消息广播到每一个跟其绑定的 queue所以也叫广播模式。publisher → Fanout exchange ├─→ queue1 → consumer1订单服务 └─→ queue2 → consumer3日志服务典型场景下单成功 → 通知订单服务更新状态 日志服务记录 积分服务加积分各业务方互不影响、各自处理。小结交换机的作用接收 publisher 发送的消息、按规则路由到与之绑定的队列FanoutExchange 会把消息路由到每个绑定的队列。7.2 Direct 交换机定向路由Direct Exchange 会将接收到的消息根据规则路由到指定的 Queue因此称为定向路由。规则有三条每一个 Queue 都与 Exchange 设置一个BindingKey发布者发送消息时指定消息的RoutingKeyExchange 将消息路由到BindingKey 与消息 RoutingKey 一致的队列。publisher → Direct exchange ├─ queue1 (BindingKey: red, blue) → consumer1 收 blue/red └─ queue2 (BindingKey: yellow, red) → consumer2 收 yellow/red发送消息时 RoutingKey 为red时queue1 和 queue2 都会收到RoutingKey 为yellow时只有 queue2 收到。7.3 Topic 交换机话题TopicExchange 与 DirectExchange 类似区别在于 routingKey 可以是多个单词的列表并且以.分割。Queue 与 Exchange 指定 BindingKey 时可以使用通配符#代表 0 个或多个单词*代表 1 个单词。比如要区分中国/日本 新闻/天气四种消息china.news 代表中国的新闻消息 china.weather 代表中国的天气消息 japan.news 代表日本新闻 japan.weather 代表日本的天气消息publisher → Topic exchange ├─ queue1 (china.#) → consumer1 收 china.* ├─ queue2 (japan.#) → consumer2 收 japan.* ├─ queue3 (#.weather) → consumer3 收 *.weather └─ queue4 (#.news) → consumer4 收 *.newsDirect 与 Topic 的差别Topic 交换机接收的消息 RoutingKey 可以是多个单词以.分隔Topic 交换机与队列绑定时的 BindingKey 可以指定通配符#代表 0 个或多个单词*代表 1 个单词。 Topic 模型是业务最常用的能用一套规则覆盖很多场景省份业务类型、商品类目操作类型…。八、声明队列和交换机SpringAMQP 提供几种方式声明队列、交换机及其绑定关系Queue声明队列可以用工厂类QueueBuilder构建Exchange声明交换机可以用ExchangeBuilder构建Binding声明队列和交换机的绑定关系可以用BindingBuilder构建。下面两种写法都要掌握。8.1 基于 Bean 的方式例如声明一个 Fanout 类型的交换机并且创建两个队列与它绑定Configuration public class FanoutConfig { // 声明 fanoutExchange 交换机 Bean public FanoutExchange fanoutExchange(){ return new FanoutExchange(hmall.fanout); } // 声明第 1 个队列 Bean public Queue fanoutQueue1(){ return new Queue(fanout.queue1); } // 绑定队列 1 和交换机 Bean public Binding bindingQueue1(Queue fanoutQueue1, FanoutExchange fanoutExchange){ return BindingBuilder.bind(fanoutQueue1).to(fanoutExchange); } // ... 略以相同方式声明第 2 个队列并完成绑定 }8.2 基于 RabbitListener 注解的方式除了写ConfigurationSpringAMQP 还提供了基于RabbitListener注解来声明队列和交换机的方式java RabbitListener(bindings QueueBinding( value Queue(name direct.queue1), exchange Exchange(name itcast.direct, type ExchangeTypes.DIRECT), key {red, blue} )) public void listenDirectQueue1(String msg){ System.out.println(消费者 1 接收到 Direct 消息【 msg 】); } 8.3 小结声明队列、交换机、绑定关系的 Bean 是什么QueueFanoutExchange / DirectExchange / TopicExchangeBinding。基于RabbitListener注解声明队列和交换机时常用的注解QueueExchange。九、消息转换器Spring 对消息对象的处理是由org.springframework.amqp.support.converter.MessageConverter负责的。默认实现是SimpleMessageConverter它基于JDK 的ObjectOutputStream完成序列化。这种默认实现有三个致命问题JDK 序列化有安全风险反序列化漏洞一直存在JDK 序列化的消息太大把类名、字段全写进去浪费带宽JDK 序列化的消息可读性差抓包看到的是二进制没法直接看。解决思路是把MessageConverter替换成更轻量、更通用的格式比如JSON。生产环境最常用的就是 Jackson2JsonMessageConverter这部分会在进阶篇中展开。十、写在最后整篇基础篇其实讲的是一件事消息从 publisher 出发经过 exchange 和 queue最终被 consumer 消费。围绕这条主线我们掌握了为什么要用异步解耦 削峰 提速怎么选型业务可靠性选 RabbitMQ海量日志选 KafkaRabbitMQ 核心组件virtual-host / publisher / consumer / queue / exchangeSpringAMQP 入门依赖、配置、发送、接收一气呵成Work 模型多消费者 prefetch 实现能者多劳三种交换机Fanout 广播、Direct 定向、Topic 通配符声明方式Bean 形式 vs 注解形式消息转换器JDK 序列化不可用要换成 JSON 之类更友好的方案。下一步进入高级篇继续解锁发送者确认、消费者确认、失败重试、延迟消息、消息堆积处理等踩坑必备技能。如果本文对你有帮助欢迎点赞 、收藏 ⭐、评论 三连。后续会更新 RabbitMQ 高级篇敬请期待。
返回列表