栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 大数据 > 大数据系统

rabbitMQ学习与使用

rabbitMQ学习与使用

RabbitMQ学习与使用
  • 1、 消息中间件概述
    • 1.1、MQ概述
    • 1.2、MQ的优势
      • 1.2.1、应用解耦
      • 1.2.2、任务异步处理
      • 1.2.3、削峰填谷
    • 1.3、MQ的劣势
    • 1.4、 常见的 MQ 产品
    • 1.5、AMQP 和 JMS
    • 1.6、 RabbitMQ
  • 2、安装与配置RabbitMQ
    • 2.1、配置虚拟主机及用户
      • 2.1.1、用户角色
      • 2.1.2、Virtual Hosts配置
        • 2.1.2.1 创建Virtual Hosts
        • 2.1.2.2、设置Virtual Hosts权限
  • 3、RabbitMQ入门
    • 3.1、搭建示例工程
      • 3.1.1、创建工程
      • 3.1.2、添加依赖
    • 3.2、编写生产者
    • 3.3、编写消费者
    • 3.4、小结
  • 4、RabbitMQ工作模式
    • 4.1、 Work queues工作队列模式
      • 4.1.1、 模式说明
      • 4.1.2、代码
      • 4.1.3、测试
      • 4.1.4、小节
    • 4.2、订阅模式概述
    • 4.3、Publish/Subscribe发布与订阅模式
      • 4.3.1、模式说明
      • 4.3.2、代码
      • 4.3.3、测试
      • 4.3.4、小结
    • 4.4、 Routing路由模式
      • 4.4.1、模式说明
      • 4.4.2、代码
      • 4.4.3、测试
      • 4.4.4、小结
    • 4.5、Topics通配符模式
      • 4.5.1、模式说明
      • 4.5.2、代码
      • 4.5.3、测试
      • 4.5.4、小结
    • 4.6、 模式总结
  • 5、Spring 整合RabbitMQ
    • 5.1、搭建生产者工程
      • 5.1.1、创建工程
      • 5.1.2、添加依赖
      • 5.1.3、配置整合
      • 5.1.4. 发送消息
    • 5.2、搭建消费者工程
      • 5.2.1、创建工程
      • 5.2.2、添加依赖
      • 5.2.3、配置整合
      • 5.2.4. 消息监听器
      • 5.2.4. 测试
  • 6. Spring Boot整合RabbitMQ
    • 6.1、搭建生产者工程
      • 6.1.1、创建工程
      • 6.1.2、添加依赖
      • 6.1.3、启动类
      • 6.1.4、配置RabbitMQ
      • 6.1.5、测试
    • 6.2、搭建消费者工程
      • 6.2.1、创建工程
      • 6.2.2、添加依赖
      • 6.2.3、启动类
      • 6.2.4、配置RabbitMQ
      • 6.2.5、消息监听处理类
      • 6.2.6、测试
  • 7、高级特性
    • 7.1、消息的可靠投递
      • 7.1.1、确认模式
      • 7.2.2、退回模式
    • 7.2、Consumer Ack
      • 7.2.1、创建consumer项目
      • 7.2.2、添加依赖
      • 7.2.3、创建配置文件
      • 7.2.4、创建监听器
      • 7.2.5、测试
    • 7.3、消费端限流
      • 7.3.1、创建并配置监听
    • 7.4、TTL
    • 7.5、死信队列
    • 7.6、延迟队列

1、 消息中间件概述 1.1、MQ概述

MQ全称 Message Queue(消息队列),是在消息的传输过程中保存消息的容器。多用于分布式系统之间进行通信。

应用之间的远程调用

加入MQ后应用之间的调用

1.2、MQ的优势 1.2.1、应用解耦

MQ相当于一个中介,生产方通过MQ与消费方交互,它将应用程序进行解耦合。

系统的耦合性越高,容错性就越低,可维护性就越低。

使用 MQ 使得应用间解耦,提升容错性和可维护性。

1.2.2、任务异步处理

将不需要同步处理的并且耗时长的操作由消息队列通知消息接收方进行异步处理。提高了应用程序的响应时间。


一个下单操作耗时:20 + 300 + 300 + 300 = 920ms
用户点击完下单按钮后,需要等待920ms才能得到下单响应,太慢!


用户点击完下单按钮后,只需等待25ms就能得到下单响应 (20 + 5 = 25ms)。
提升用户体验和系统吞吐量(单位时间内处理请求的数目)。

1.2.3、削峰填谷

如订单系统,在下单的时候就会往数据库写数据。但是数据库只能支撑每秒1000左右的并发写入,并发量再高就容易宕机。低峰期的时候并发也就100多个,但是在高峰期时候,并发量会突然激增到5000以上,这个时候数据库肯定卡死了。

消息被MQ保存起来了,然后系统就可以按照自己的消费能力来消费,比如每秒1000个消息,这样慢慢写入数据库,这样就不会卡死数据库了。


但是使用了MQ之后,限制消费消息的速度为1000,但是这样一来,高峰期产生的数据势必会被积压在MQ中,高峰就被“削”掉了。但是因为消息积压,在高峰期过后的一段时间内,消费消息的速度还是会维持在1000QPS,直到消费完积压的消息,这就叫做“填谷”

1.3、MQ的劣势

系统可用性降低
系统引入的外部依赖越多,系统稳定性越差。一旦 MQ 宕机,就会对业务造成影响。如何保证MQ的高可用?

系统复杂度提高
MQ 的加入大大增加了系统的复杂度,以前系统间是同步的远程调用,现在是通过 MQ 进行异步调用。如何保证消息没有被重复消费?怎么处理消息丢失情况?那么保证消息传递的顺序性?

一致性问题
A 系统处理完业务,通过 MQ 给B、C、D三个系统发消息,如果 B 系统、C 系统处理成功,D 系统处理失败。如何保证消息数据处理的一致性?

1.4、 常见的 MQ 产品

目前业界有很多的 MQ 产品,例如 RabbitMQ、RocketMQ、ActiveMQ、Kafka、ZeroMQ、metaMq等,也有直接使用 Redis 充当消息队列的案例,而这些消息队列产品,各有侧重,在实际选型时,需要结合自身需求及 MQ 产品特征,综合考虑。

RabbitMQActiveMQRocketMQKafka
公司/ 社区RabbitApache阿里Apache
开发语言ErlangJavaJavaScala&Java
协议支持AMQP,XMPP,SMTP,STOMPOpenWire,STOMP,REST,XMPP,AMQP自定义自定义协议,社区封装了http协议支持
客户端支持语言官方支持Erlang,Java,Ruby等,社区产出多种API,几乎支持所有语言Java,C,C++,Python,PHP,Perl,.net等Java,C++(不成熟)官方支持Java,社区产出多种API,如PHP,Python等
单机吞吐量万级(其次)万级(最差)十万级(最好)十万级(次之)
消息延迟微妙级毫秒级毫秒级毫秒以内
功能特性并发能力强,性能极其好,延时低,社区活跃,管理界面丰富老牌产品,成熟度高,文档较多MQ功能比较完备,扩展性佳只支持主要的MQ功能,毕竟是为大数据领域准备的。
1.5、AMQP 和 JMS

实现MQ的大致有两种主流方式:AMQP、JMS。
AMQP
AMQP,即 Advanced Message Queuing Protocol(高级消息队列协议),是一个网络协议,是应用层协议的一个开放标准,为面向消息的中间件设计。基于此协议的客户端与消息中间件可传递消息,遵循此协议,不收客户端和中间件产品和开发语言限制。2006年,AMQP 规范发布。类比HTTP。

JMS
JMS 即 Java 消息服务(JavaMessage Service)应用程序接口,是一个 Java 平台中关于面向消息中间件的API

JMS 是 JavaEE 规范中的一种,类比JDBC

很多消息中间件都实现了JMS规范,例如:ActiveMQ。RabbitMQ 官方没有提供 JMS 的实现包,但是开源社区有

AMQP 与 JMS 区别

  • JMS是定义了统一的接口,来对消息操作进行统一;AMQP是通过规定协议来统一数据交互的格式
  • JMS限定了必须使用Java语言;AMQP只是协议,不规定实现方式,因此是跨语言的。
  • JMS规定了两种消息模式;而AMQP的消息模式更加丰富
1.6、 RabbitMQ

RabbitMQ官方地址:http://www.rabbitmq.com/

2007年,Rabbit 技术公司基于 AMQP 标准开发的 RabbitMQ 1.0 发布。RabbitMQ 采用 Erlang 语言开
发。Erlang 语言专门为开发高并发和分布式系统的一种语言,在电信领域使用广泛。
RabbitMQ 基础架构如下图:

RabbitMQ 中的相关概念:

  • Broker:接收和分发消息的应用,RabbitMQ Server就是 Message Broker
  • Virtual host:出于多租户和安全因素设计的,把 AMQP 的基本组件划分到一个虚拟的分组中,类似于网络中的 namespace 概念。当多个不同的用户使用同一个 RabbitMQ server 提供的服务时,可以划分出多个vhost,每个用户在自己的 vhost 创建 exchange/queue 等
  • Connection:publisher/consumer 和 broker 之间的 TCP 连接
  • Channel:如果每一次访问 RabbitMQ 都建立一个 Connection,在消息量大的时候建立 TCP
    Connection的开销将是巨大的,效率也较低。Channel 是在 connection 内部建立的逻辑连接,如果应用程序支持多线程,通常每个thread创建单独的 channel 进行通讯,AMQP method 包含了channel id 帮助客户端和message broker 识别 channel,所以 channel 之间是完全隔离的。Channel 作为轻量级的 Connection 极大减少了操作系统建立 TCP connection 的开销
  • Exchange:message 到达 broker 的第一站,根据分发规则,匹配查询表中的 routing key,分发消息到queue 中去。常用的类型有:direct (point-to-point), topic (publish-subscribe) and fanout (multicast)
  • Queue:消息最终被送到这里等待 consumer 取走
  • Binding:exchange 和 queue 之间的虚拟连接,binding 中可以包含routing key。Binding 信息被保存到 exchange 中的查询表中,用于 message 的分发依据

RabbitMQ提供了6种模式:简单模式,work模式,Publish/Subscribe发布与订阅模式,Routing路由模式,Topics主题模式,RPC远程调用模式(远程调用,不太算MQ;暂不作介绍);
官网对应模式介绍:https://www.rabbitmq.com/getstarted.html

2、安装与配置RabbitMQ

安装教程请移步

2.1、配置虚拟主机及用户 2.1.1、用户角色

RabbitMQ在安装号后,可以访问 http://ip:15672,用默认的guest/guest用户密码进行登陆,创建用户也可以用管控台添加,如下图:


角色说明:

  1. 超级管理员(administrator)
    可登陆管理控制台,可查看所有的信息,并且可以对⽤户,策略(policy)进⾏操作。
  2. 监控者(monitoring)
    可登陆管理控制台,同时可以查看rabbitmq节点的相关信息(进程数,内存使⽤情况,磁盘使⽤情况等)
  3. 策略制定者(policymaker)
    可登陆管理控制台, 同时可以对policy进⾏管理。但⽆法查看节点的相关信息(上图红框标识的部分)。
  4. 普通管理者(management)
    仅可登陆管理控制台,⽆法看到节点信息,也⽆法对策略进⾏管理。
  5. 其他
    ⽆法登陆管理控制台,通常就是普通的⽣产者和消费者。
2.1.2、Virtual Hosts配置

像mysql拥有数据库的概念并且可以指定⽤户对库和表等操作的权限。RabbitMQ也有类似的权限管理;在RabbitMQ中可以虚拟消息服务器Virtual Host,每个Virtual Hosts相当于⼀个相对独⽴的RabbitMQ服务器,每个VirtualHost之间是相互隔离的。exchange、queue、message不能互通。 相当于mysql的db。Virtual Name⼀般以/开头。

2.1.2.1 创建Virtual Hosts

2.1.2.2、设置Virtual Hosts权限

点击Virtual Hosts的名称

展开Permissions,选择设置当前Virtual Hosts的拥有者

3、RabbitMQ入门

简单模式

在上图的模型中,有以下概念:

  • P:生产者,也就是要发送消息的程序
  • C:消费者:消息的接收者,会一直等待消息到来
  • queue:消息队列,图中红色部分。类似一个邮箱,可以缓存消息;生产者向其中投递消息,消费者从其中取出消息
3.1、搭建示例工程 3.1.1、创建工程 3.1.2、添加依赖

    com.rabbitmq
    amqp-client
    5.13.1

3.2、编写生产者

编写消息生产者com.flaw.rabbitmq.simple.Producer

package com.flaw.rabbitmq.simple;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class Producer {

    static final String QUEUE_NAME = "simple_queue";

    public static void main(String[] args) throws Exception {
        // 创建连接工厂
        ConnectionFactory connectionFactory=new ConnectionFactory();
        // 主机地址;默认为 localhost
        connectionFactory.setHost("localhost");
        // 连接端口;默认为 5672
        connectionFactory.setPort(5672);
        // 虚拟主机名称;默认为 /
        connectionFactory.setVirtualHost("/xzk");
        // 连接用户名;默认为 guest
        connectionFactory.setUsername("flaw");
        // 连接用户名;默认为 guest
        connectionFactory.setPassword("flaw");
        // 创建连接
        Connection connection=connectionFactory.newConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明(创建)队列
        
        channel.queueDeclare(QUEUE_NAME,true,false,false,null);
        // 要发送的信息
        String message="hello,RabbitMQ!";
        
        channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
        System.out.println("以发送消息:"+message);
        // 关闭释放资源
        channel.close();
        connection.close();
    }

}

在执行上述的消息发送之后;可以登录rabbitMQ的管理控制台,可以发现队列和其消息:

3.3、编写消费者

抽取创建connection的工具类com.flaw.rabbitmq.util.ConnectionUtil;

package com.flaw.rabbitmq.util;

import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class ConnectionUtil {

    public static Connection getConnection() throws Exception {
        // 创建连接工厂
        ConnectionFactory connectionFactory=new ConnectionFactory();
        // 主机地址;默认为 localhost
        connectionFactory.setHost("localhost");
        // 连接端口;默认为 5672
        connectionFactory.setPort(5672);
        // 虚拟主机名称;默认为 /
        connectionFactory.setVirtualHost("/xzk");
        // 连接用户名;默认为 guest
        connectionFactory.setUsername("flaw");
        // 连接用户名;默认为 guest
        connectionFactory.setPassword("flaw");
        //创建连接返回
        return connectionFactory.newConnection();
    }

}

编写消息的消费者com.flaw.rabbitmq.simple.Consumer

package com.flaw.rabbitmq.simple;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明(创建)队列
        
        channel.queueDeclare(Producer.QUEUE_NAME,true,false,false,null);
        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                //路由key
                System.out.println("路由key为:" + envelope.getRoutingKey());
                //交换机
                System.out.println("交换机为:" + envelope.getExchange());
                //消息id
                System.out.println("消息id为:" + envelope.getDeliveryTag());
                //收到的消息
                System.out.println("接收到的消息为:" + new String(body, "utf-8"));
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.QUEUE_NAME, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

3.4、小结

上述的入门案例中中其实使用的是如下的简单模式:

在上图的模型中,有以下概念:

  • P:生产者,也就是要发送消息的程序
  • C:消费者:消息的接受者,会一直等待消息到来。
  • queue:消息队列,图中红色部分。类似一个邮箱,可以缓存消息;生产者向其中投递消息,消费者从其中取出消息。
4、RabbitMQ工作模式 4.1、 Work queues工作队列模式 4.1.1、 模式说明


Work Queues 与入门程序的 简单模式 相比,多了一个或一些消费端,多个消费端共同消费同一个队列中的消息。
应用场景: 对于 任务过重或任务较多情况使用工作队列可以提高任务处理的速度。

4.1.2、代码

Work Queues 与入门程序的 简单模式 的代码是几乎一样的;可以完全复制,并复制多一个消费者进行多个消费者同时消费消息的测试。
生产者

package com.flaw.rabbitmq.work;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class Producer {

    static final String QUEUE_NAME = "work_queue";

    public static void main(String[] args) throws Exception {
        // 创建连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明(创建)队列
        
        channel.queueDeclare(QUEUE_NAME,true,false,false,null);
        for (int i = 1; i <= 30; i++) {
            // 要发送的信息
            String message="hello,RabbitMQ! work模式==="+i;
            
            channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
            System.out.println("已发送消息:"+message);
        }
        // 关闭释放资源
        channel.close();
        connection.close();
    }

}

消费者1

package com.flaw.rabbitmq.work;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer1 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明(创建)队列
        
        channel.queueDeclare(Producer.QUEUE_NAME,true,false,false,null);

        // 一次只能接收并处理一个消息
        channel.basicQos(1);

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                try {
                    //路由key
                    System.out.println("路由key为:" + envelope.getRoutingKey());
                    //交换机
                    System.out.println("交换机为:" + envelope.getExchange());
                    //消息id
                    System.out.println("消息id为:" + envelope.getDeliveryTag());
                    //收到的消息
                    System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8"));

                    // 休眠1秒,测试时观察更直观
                    Thread.sleep(1000);
                    //确认消息
                    channel.basicAck(envelope.getDeliveryTag(), false);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.QUEUE_NAME, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

消费者2

package com.flaw.rabbitmq.work;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer2 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明(创建)队列
        
        channel.queueDeclare(Producer.QUEUE_NAME,true,false,false,null);

        // 一次只能接收并处理一个消息
        channel.basicQos(1);

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                try {
                    //路由key
                    System.out.println("路由key为:" + envelope.getRoutingKey());
                    //交换机
                    System.out.println("交换机为:" + envelope.getExchange());
                    //消息id
                    System.out.println("消息id为:" + envelope.getDeliveryTag());
                    //收到的消息
                    System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8"));

                    // 休眠1秒,测试时观察更直观
                    Thread.sleep(1000);
                    //确认消息
                    channel.basicAck(envelope.getDeliveryTag(), false);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.QUEUE_NAME, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

4.1.3、测试

启动两个消费者,然后再启动生产者发送消息;到IDEA的两个消费者对应的控制台查看是否竞争性的接收到消息。


4.1.4、小节

在一个队列中如果有多个消费者,那么消费者之间对于同一个消息的关系是竞争的关系。

4.2、订阅模式概述

订阅模式示例图:

前面2个案例中,只有3个角色:

  • P:生产者,也就是要发送消息的程序
  • C:消费者:消息的接受者,会一直等待消息到来。
  • queue:消息队列,图中红色部分

而在订阅模型中,多了一个exchange角色,而且过程略有变化:

  • P:生产者,也就是要发送消息的程序,但是不再发送到队列中,而是发给X(交换机)
  • C:消费者,消息的接受者,会一直等待消息到来。
  • Queue:消息队列,接收消息、缓存消息。
  • Exchange:交换机,图中的X。一方面,接收生产者发送的消息。另一方面,知道如何处理消息,例如递交给某个特别队列、递交给所有队列、或是将消息丢弃。到底如何操作,取决于Exchange的类型。Exchange有常见以下3种类型:
    • Fanout:广播,将消息交给所有绑定到交换机的队列
    • Direct:定向,把消息交给符合指定routing key 的队列
    • Topic:通配符,把消息交给符合routing pattern(路由模式) 的队列

Exchange(交换机)只负责转发消息,不具备存储消息的能力,因此如果没有任何队列与Exchange绑定,或者没有符合路由规则的队列,那么消息会丢失!

4.3、Publish/Subscribe发布与订阅模式 4.3.1、模式说明


发布订阅模式:

  1. 每个消费者监听自己的队列。
  2. 生产者将消息发给broker,由交换机将消息转发到绑定此交换机的每个队列,每个绑定交换机的队列都将接收到消息
4.3.2、代码

生产者:

package com.flaw.rabbitmq.ps;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;


public class Producer {

    
    static final String FANOUT_EXCHANGE = "fanout_exchange";
    
    static final String FANOUT_QUEUE_1 = "fanout_queue_1";
    
    static final String FANOUT_QUEUE_2 = "fanout_queue_2";

    public static void main(String[] args) throws Exception {
        // 创建连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        
        channel.exchangeDeclare(FANOUT_EXCHANGE, BuiltinExchangeType.FANOUT);
        
        channel.queueDeclare(FANOUT_QUEUE_1,true,false,false,null);
        channel.queueDeclare(FANOUT_QUEUE_2,true,false,false,null);
        
        channel.queueBind(FANOUT_QUEUE_1, FANOUT_EXCHANGE, "");
        channel.queueBind(FANOUT_QUEUE_2, FANOUT_EXCHANGE, "");

        for (int i = 1; i <= 10; i++) {
            // 要发送的信息
            String message="hello,RabbitMQ! 发布订阅模式==="+i;
            
            channel.basicPublish(FANOUT_EXCHANGE,"",null,message.getBytes());
            System.out.println("已发送消息:"+message);
        }
        // 关闭释放资源
        channel.close();
        connection.close();
    }

}

消费者1:

package com.flaw.rabbitmq.ps;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer1 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明交换机
        channel.exchangeDeclare(Producer.FANOUT_EXCHANGE, BuiltinExchangeType.FANOUT);
        
        channel.queueDeclare(Producer.FANOUT_QUEUE_1,true,false,false,null);
        // 队列绑定交换机
        channel.queueBind(Producer.FANOUT_QUEUE_1, Producer.FANOUT_EXCHANGE, "");

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                //路由key
                System.out.println("路由key为:" + envelope.getRoutingKey());
                //交换机
                System.out.println("交换机为:" + envelope.getExchange());
                //消息id
                System.out.println("消息id为:" + envelope.getDeliveryTag());
                //收到的消息
                System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8"));
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.FANOUT_QUEUE_1, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

消费者2:

package com.flaw.rabbitmq.ps;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer2 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明交换机
        channel.exchangeDeclare(Producer.FANOUT_EXCHANGE, BuiltinExchangeType.FANOUT);
        
        channel.queueDeclare(Producer.FANOUT_QUEUE_2,true,false,false,null);
        // 队列绑定交换机
        channel.queueBind(Producer.FANOUT_QUEUE_2, Producer.FANOUT_EXCHANGE, "");

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                //路由key
                System.out.println("路由key为:" + envelope.getRoutingKey());
                //交换机
                System.out.println("交换机为:" + envelope.getExchange());
                //消息id
                System.out.println("消息id为:" + envelope.getDeliveryTag());
                //收到的消息
                System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8"));
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.FANOUT_QUEUE_2, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

4.3.3、测试

启动所有消费者,然后使用生产者发送消息;在每个消费者对应的控制台可以查看到生产者发送的所有消息;到达广播(fanout模式)的效果。

在执行完测试代码后,其实到RabbitMQ的管理后台找到 Exchanges 选项卡,点击 fanout_exchange的交换机,可以查看到如下的绑定:

4.3.4、小结

交换机需要与队列进行绑定,绑定之后;一个消息可以被多个消费者都收到。

发布订阅模式与工作队列模式的区别:

  1. 工作队列模式不用定义交换机,而发布/订阅模式需要定义交换机。
  2. 发布/订阅模式的生产方是面向交换机发送消息,工作队列模式的生产方是面向队列发送消息(底层使用默认交换机)。
  3. 发布/订阅模式需要设置队列和交换机的绑定,工作队列模式不需要设置,实际上工作队列模式会将队列绑 定到默认的交换机 。
4.4、 Routing路由模式 4.4.1、模式说明

路由模式特点:

  • 队列与交换机的绑定,不能是任意绑定了,而是要指定一个 RoutingKey(路由key)
  • 消息的发送方在 向 Exchange发送消息时,也必须指定消息的 RoutingKey。
  • Exchange不再把消息交给每一个绑定的队列,而是根据消息的 RoutingKey 进行判断,只有队列的Routingkey与消息的 Routingkey 完全一致,才会接收到消息

    图解:
  • P:生产者,向Exchange发送消息,发送消息时,会指定一个routingkey。
  • X:Exchange(交换机),接收生产者的消息,然后把消息递交给与routingkey完全匹配的队列
  • C1:消费者,其所在队列指定了需要routingkey 为 error 的消息
  • C2:消费者,其所在队列指定了需要routingkey 为 info、error、warning 的消息
4.4.2、代码

在编码上与 Publish/Subscribe发布与订阅模式的区别是交换机的类型为:Direct,还有队列绑定交换机的时候需要指定routingkey。

生产者:

package com.flaw.rabbitmq.routing;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;


public class Producer {

    
    static final String DIRECT_EXCHANGE = "direct_exchange";
    
    static final String DIRECT_QUEUE_INSERT = "direct_queue_insert";
    
    static final String DIRECT_QUEUE_UPDATE = "direct_queue_update";

    public static void main(String[] args) throws Exception {
        // 创建连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        
        channel.exchangeDeclare(DIRECT_EXCHANGE, BuiltinExchangeType.DIRECT);
        
        channel.queueDeclare(DIRECT_QUEUE_INSERT,true,false,false,null);
        channel.queueDeclare(DIRECT_QUEUE_UPDATE,true,false,false,null);
        
        channel.queueBind(DIRECT_QUEUE_INSERT, DIRECT_EXCHANGE, "insert");
        channel.queueBind(DIRECT_QUEUE_UPDATE, DIRECT_EXCHANGE, "update");

        // 要发送的信息
        String message="新增了商品! 路由模式:routingKey为 insert";
        
        channel.basicPublish(DIRECT_EXCHANGE,"insert",null,message.getBytes());
        System.out.println("已发送消息:"+message);

        // 要发送的信息
        message="修改了商品! 路由模式:routingKey为 update";
        
        channel.basicPublish(DIRECT_EXCHANGE,"update",null,message.getBytes());
        System.out.println("已发送消息:"+message);
        // 关闭释放资源
        channel.close();
        connection.close();
    }

}

消费者1:

package com.flaw.rabbitmq.routing;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer1 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明交换机
        channel.exchangeDeclare(Producer.DIRECT_EXCHANGE, BuiltinExchangeType.DIRECT);
        
        channel.queueDeclare(Producer.DIRECT_QUEUE_INSERT,true,false,false,null);
        // 队列绑定交换机
        channel.queueBind(Producer.DIRECT_QUEUE_INSERT, Producer.DIRECT_EXCHANGE, "insert");

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                //路由key
                System.out.println("路由key为:" + envelope.getRoutingKey());
                //交换机
                System.out.println("交换机为:" + envelope.getExchange());
                //消息id
                System.out.println("消息id为:" + envelope.getDeliveryTag());
                //收到的消息
                System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8"));
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.DIRECT_QUEUE_INSERT, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

消费者2:

package com.flaw.rabbitmq.routing;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer2 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明交换机
        channel.exchangeDeclare(Producer.DIRECT_EXCHANGE, BuiltinExchangeType.DIRECT);
        
        channel.queueDeclare(Producer.DIRECT_QUEUE_UPDATE,true,false,false,null);
        // 队列绑定交换机
        channel.queueBind(Producer.DIRECT_QUEUE_UPDATE, Producer.DIRECT_EXCHANGE, "update");

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                //路由key
                System.out.println("路由key为:" + envelope.getRoutingKey());
                //交换机
                System.out.println("交换机为:" + envelope.getExchange());
                //消息id
                System.out.println("消息id为:" + envelope.getDeliveryTag());
                //收到的消息
                System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8"));
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.DIRECT_QUEUE_UPDATE, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

4.4.3、测试

启动所有消费者,然后使用生产者发送消息;在消费者对应的控制台可以查看到生产者发送对应routingKey对应队列的消息;到达按照需要接收的效果。

在执行完测试代码后,其实到RabbitMQ的管理后台找到 Exchanges 选项卡,点击 direct_exchange的交换机,可以查看到如下的绑定:

4.4.4、小结

Routing模式要求队列在绑定交换机时要指定routingKey,消息会转发到符合routingKey的队列。

4.5、Topics通配符模式 4.5.1、模式说明

Topic 类型与 Direct 相比,都是可以根据 RoutingKey 把消息路由到不同的队列。只不过 Topic 类型Exchange 可以让队列在绑定 RoutingKey 的时候使用通配符!
Routingkey 一般都是有一个或多个单词组成,多个单词之间以"."分割,例如: item.insert
通配符规则:
# :匹配一个或多个词
*:匹配不多不少恰好1个词
举例:
item.#:能够匹配 item.insert.abc 或者 item.insert
item.* :只能匹配 item.insert


图解:

  • 红色Queue:绑定的是 usa.#,因此凡是以 usa.开头的 routingKey 都会被匹配到
  • 黄色Queue:绑定的是 #.news ,因此凡是以 .news 结尾的 routingKey 都会被匹配
4.5.2、代码


生产者:
使用topic类型的Exchange,发送消息的routingKey有3种: goods.insert 、 goods.update 、 goods.delete :

package com.flaw.rabbitmq.topics;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;


public class Producer {

    
    static final String TOPIC_EXCHANGE = "topic_exchange";
    
    static final String TOPIC_QUEUE_1 = "topic_queue_1";
    
    static final String TOPIC_QUEUE_2 = "topic_queue_2";

    public static void main(String[] args) throws Exception {
        // 创建连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        
        channel.exchangeDeclare(TOPIC_EXCHANGE, BuiltinExchangeType.TOPIC);

        // 要发送的信息
        String message="新增了商品! Topic模式:routingKey为 goods.insert";
        
        channel.basicPublish(TOPIC_EXCHANGE,"goods.insert",null,message.getBytes());
        System.out.println("已发送消息:"+message);

        // 要发送的信息
        message="修改了商品! Topic模式:routingKey为 goods.update";
        
        channel.basicPublish(TOPIC_EXCHANGE,"goods.update",null,message.getBytes());
        System.out.println("已发送消息:"+message);

        // 要发送的信息
        message="删除了商品! Topic模式:routingKey为 goods.delete";
        
        channel.basicPublish(TOPIC_EXCHANGE,"goods.delete",null,message.getBytes());
        System.out.println("已发送消息:"+message);
        // 关闭释放资源
        channel.close();
        connection.close();
    }

}

消费者1:
接收两种类型的消息:更新商品goods.update和删除商品goods.delete

package com.flaw.rabbitmq.topics;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer1 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        //声明交换机
        channel.exchangeDeclare(Producer.TOPIC_EXCHANGE, BuiltinExchangeType.TOPIC);
        
        channel.queueDeclare(Producer.TOPIC_QUEUE_1,true,false,false,null);
        // 队列绑定交换机
        channel.queueBind(Producer.TOPIC_QUEUE_1, Producer.TOPIC_EXCHANGE, "goods.update");
        channel.queueBind(Producer.TOPIC_QUEUE_1, Producer.TOPIC_EXCHANGE, "goods.delete");

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                //路由key
                System.out.println("路由key为:" + envelope.getRoutingKey());
                //交换机
                System.out.println("交换机为:" + envelope.getExchange());
                //消息id
                System.out.println("消息id为:" + envelope.getDeliveryTag());
                //收到的消息
                System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8"));
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.TOPIC_QUEUE_1, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

消费者2:
接收所有类型的消息:新增商品,更新商品和删除商品(goods.*)。

package com.flaw.rabbitmq.topics;

import com.flaw.rabbitmq.util.ConnectionUtil;
import com.rabbitmq.client.*;

import java.io.IOException;

public class Consumer2 {
    public static void main(String[] args) throws Exception {
        // 获取连接
        Connection connection= ConnectionUtil.getConnection();
        // 创建频道
        Channel channel = connection.createChannel();
        // 声明交换机
        channel.exchangeDeclare(Producer.TOPIC_EXCHANGE, BuiltinExchangeType.TOPIC);
        
        channel.queueDeclare(Producer.TOPIC_QUEUE_2,true,false,false,null);
        // 队列绑定交换机
        channel.queueBind(Producer.TOPIC_QUEUE_2, Producer.TOPIC_EXCHANGE, "goods.*");

        // 创建消费者,并设置消息处理
        DefaultConsumer consumer=new DefaultConsumer(channel){
            
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                //路由key
                System.out.println("路由key为:" + envelope.getRoutingKey());
                //交换机
                System.out.println("交换机为:" + envelope.getExchange());
                //消息id
                System.out.println("消息id为:" + envelope.getDeliveryTag());
                //收到的消息
                System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8"));
            }
        };
        // 监听消息
        
        channel.basicConsume(Producer.TOPIC_QUEUE_2, true, consumer);
        //需要一直监听消息,所以不需要关闭释放资源
    }
}

4.5.3、测试

启动所有消费者,然后使用生产者发送消息;在消费者对应的控制台可以查看到生产者发送对应routingKey对应队列的消息;到达按照需要接收的效果;并且这些routingKey可以使用通配符。

在执行完测试代码后,其实到RabbitMQ的管理后台找到 Exchanges 选项卡,点击 topic_exchange的交换机,可以查看到如下的绑定:

4.5.4、小结

Topic主题模式可以实现 Publish/Subscribe发布与订阅模式 和 Routing路由模式 的功能;只是Topic在配置routingKey 的时候可以使用通配符,显得更加灵活。

4.6、 模式总结

RabbitMQ工作模式:

  1. 简单模式 HelloWorld 一个生产者、一个消费者,不需要设置交换机(使用默认的交换机)
  2. 工作队列模式 Work Queue 一个生产者、多个消费者(竞争关系),不需要设置交换机(使用默认的交换机)
  3. 发布订阅模式 Publish/subscribe 需要设置类型为fanout的交换机,并且交换机和队列进行绑定,当发送消息到交换机后,交换机会将消息发送到绑定的队列
  4. 路由模式 Routing 需要设置类型为direct的交换机,交换机和队列进行绑定,并且指定routingKey,当发送消息到交换机后,交换机会根据routing key将消息发送到对应的队列
  5. 通配符模式 Topic 需要设置类型为topic的交换机,交换机和队列进行绑定,并且指定通配符方式的routingKey,当发送消息到交换机后,交换机会根据routing key将消息发送到对应的队列
5、Spring 整合RabbitMQ 5.1、搭建生产者工程 5.1.1、创建工程

5.1.2、添加依赖

修改pom.xml文件内容为如下:



    4.0.0

    com.flaw
    spring-rabbitmq-producer
    1.0-SNAPSHOT

    
        
            org.springframework
            spring-context
            5.1.7.RELEASE
        
        
            org.springframework.amqp
            spring-rabbit
            2.1.8.RELEASE
        
        
            junit
            junit
            4.12
        
        
            org.springframework
            spring-test
            5.1.7.RELEASE
        
    

5.1.3、配置整合
  1. 创建 spring-rabbitmq-producersrcmainresourcespropertiesrabbitmq.properties连接参数等配置文件
rabbitmq.host=192.168.56.1
rabbitmq.port=5672
rabbitmq.username=flaw
rabbitmq.password=flaw
rabbitmq.virtual-host=/xzk
  1. 创建 spring-rabbitmq-producersrcmainresourcesspringspring-rabbitmq.xml 整合配置文件


    
    
    
    
    
    
    
    
    
    
    
    
    
    
    
        
            
            
        
    
    
    
    
    
    
    
    
    
        
            
            
            
        
    
    
    

5.1.4. 发送消息

创建测试文件 spring-rabbitmq-producersrctestjavacomflawrabbitmqProducerTest.java

package com.flaw.rabbitmq;

import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;

@RunWith(SpringJUnit4ClassRunner.class) 
@ContextConfiguration(locations = "classpath:spring/spring-rabbitmq.xml")
public class ProducerTest {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    
    @Test
    public void queueTest(){
        // 路由键与队列同名
        rabbitTemplate.convertAndSend("spring_queue","发送队列spring_queue的消息。");
    }

    
    @Test
    public void fanoutTest(){
        
        rabbitTemplate.convertAndSend("spring_fanout_exchange","","发送到spring_fanout_exchange交换机的广播消息");
    }

    
    @Test
    public void topicTest(){
        
        rabbitTemplate.convertAndSend("spring_topic_exchange","flaw.haha","发送到spring_topic_exchange交换机flaw.haha的消息");
        rabbitTemplate.convertAndSend("spring_topic_exchange","flaw.haha.1","发送到spring_topic_exchange交换机flaw.haha.1的消息");
        rabbitTemplate.convertAndSend("spring_topic_exchange","flaw.haha.2","发送到spring_topic_exchange交换机flaw.haha.2的消息");
        rabbitTemplate.convertAndSend("spring_topic_exchange","xzk.com","发送到spring_topic_exchange交换机xzk.com的消息");
    }
}
5.2、搭建消费者工程 5.2.1、创建工程

5.2.2、添加依赖

修改pom.xml文件内容为如下:



    4.0.0

    com.flaw
    spring-rabbitmq-consumer
    1.0-SNAPSHOT

    
        
            org.springframework
            spring-context
            5.1.7.RELEASE
        
        
            org.springframework.amqp
            spring-rabbit
            2.1.8.RELEASE
        
        
            junit
            junit
            4.12
        
        
            org.springframework
            spring-test
            5.1.7.RELEASE
        
    

5.2.3、配置整合
  1. 创建 spring-rabbitmq-consumersrcmainresourcespropertiesrabbitmq.properties连接参数等配置文件
rabbitmq.host=192.168.56.1
rabbitmq.port=5672
rabbitmq.username=flaw
rabbitmq.password=flaw
rabbitmq.virtual-host=/xzk
  1. 创建 spring-rabbitmq-consumersrcmainresourcesspringspring-rabbitmq.xml 整合配置文件
在这里插入代码片
5.2.4. 消息监听器
  1. 队列监听器
    创建 spring-rabbitmq-consumersrcmainjavacomflawrabbitmqlistenerSpringQueueListener.java监听spring_queue队列的消息
package com.flaw.rabbitmq.listener;

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageListener;

public class SpringQueueListener implements MessageListener {
    public void onMessage(Message message) {
        try {
            String msg = new String(message.getBody(), "utf-8");
            System.out.printf("接收路由名称为:%s,路由键为:%s,队列名为:%s的消息:%s n",
                    message.getMessageProperties().getReceivedExchange(),
                    message.getMessageProperties().getReceivedRoutingKey(),
                    message.getMessageProperties().getConsumerQueue(),
                    msg);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

  1. 广播监听器1
    创建 spring-rabbitmq-consumersrcmainjavacomflawrabbitmqlistenerFanoutListener1.java 监听spring_fanout_queue_1队列的消息
package com.flaw.rabbitmq.listener;

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageListener;

public class FanoutListener1 implements MessageListener {
    @Override
    public void onMessage(Message message) {
        try {
            String msg = new String(message.getBody(), "utf-8");
            System.out.printf("广播监听器1:接收路由名称为:%s,路由键为:%s,队列名为:%s的消息为:%s n",
                    message.getMessageProperties().getReceivedExchange(),
                    message.getMessageProperties().getReceivedRoutingKey(),
                    message.getMessageProperties().getConsumerQueue(),
                    msg);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

  1. 广播监听器2
    创建 spring-rabbitmq-consumersrcmainjavacomflawrabbitmqlistenerFanoutListener2.java 监听spring_fanout_queue_2的消息
package com.flaw.rabbitmq.listener;

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageListener;

public class FanoutListener2 implements MessageListener {
    @Override
    public void onMessage(Message message) {
        try {
            String msg = new String(message.getBody(), "utf-8");
            System.out.printf("广播监听器2:接收路由名称为:%s,路由键为:%s,队列名为: %s的消息:%s n",
                    message.getMessageProperties().getReceivedExchange(),
                    message.getMessageProperties().getReceivedRoutingKey(),
                    message.getMessageProperties().getConsumerQueue(),
                    msg);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}
  1. *通配符监听器
    创建 spring-rabbitmq-consumersrcmainjavacomflawrabbitmqlistenerTopicListenerStar.java 监听spring_topic_queue_star队列的消息
package com.flaw.rabbitmq.listener;

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageListener;

public class TopicListenerStar implements MessageListener {
    @Override
    public void onMessage(Message message) {
        try {
            String msg = new String(message.getBody(), "utf-8");
            System.out.printf("通配符*监听器:接收路由名称为:%s,路由键为:%s,队列名 为:%s的消息:%s n",
                    message.getMessageProperties().getReceivedExchange(),
                    message.getMessageProperties().getReceivedRoutingKey(),
                    message.getMessageProperties().getConsumerQueue(),
                    msg);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

  1. #通配符监听器1
    创建 spring-rabbitmq-consumersrcmainjavacomflawrabbitmqlistenerTopicListenerWell1.java spring_topic_queue_well队列的消息
package com.flaw.rabbitmq.listener;

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageListener;

public class TopicListenerWell1 implements MessageListener {
    @Override
    public void onMessage(Message message) {
        try {
            String msg = new String(message.getBody(), "utf-8");
            System.out.printf("通配符#监听器:接收路由名称为:%s,路由键为:%s,队列名 为:%s的消息:%s n",
                    message.getMessageProperties().getReceivedExchange(),
                    message.getMessageProperties().getReceivedRoutingKey(),
                    message.getMessageProperties().getConsumerQueue(),
                    msg);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

  1. #通配符监听器2
    创建 spring-rabbitmq-consumersrcmainjavacomflawrabbitmqlistenerTopicListenerWell2.java 监听spring_topic_queue_well2队列的消息
package com.flaw.rabbitmq.listener;

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageListener;

public class TopicListenerWell2 implements MessageListener {
    @Override
    public void onMessage(Message message) {
        try {
            String msg = new String(message.getBody(), "utf-8");
            System.out.printf("通配符#监听器2:接收路由名称为:%s,路由键为:%s,队列名 为:%s的消息:%s n",
                    message.getMessageProperties().getReceivedExchange(),
                    message.getMessageProperties().getReceivedRoutingKey(),
                    message.getMessageProperties().getConsumerQueue(),
                    msg);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

5.2.4. 测试

创建测试类:

package com.flaw.rabbitmq;

import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;

@RunWith(SpringJUnit4ClassRunner.class)
@ContextConfiguration(locations = "classpath:spring/spring-rabbitmq.xml")
public class ConsumerTest {

    @Test
    public void test(){
        while (true){
            // 死循环让consumer保持连接状态,一直监听队列里面的消息
        }
    }
}


6. Spring Boot整合RabbitMQ

在Spring项目中,可以使用Spring-Rabbit去操作RabbitMQ https://github.com/spring-projects/spring-amqp
尤其是在spring boot项目中只需要引入对应的amqp启动器依赖即可,方便的使用RabbitTemplate发送消息,使用注解接收消息。
一般在开发过程中:

生产者工程:

  1. application.yml文件配置RabbitMQ相关信息;
  2. 在生产者工程中编写配置类,用于创建交换机和队列,并进行绑定
  3. 注入RabbitTemplate对象,通过RabbitTemplate对象发送消息到交换机

消费者工程:

  1. application.yml文件配置RabbitMQ相关信息
  2. 创建消息处理类,用于接收队列中的消息并进行处理
6.1、搭建生产者工程 6.1.1、创建工程

创建生产者工程springboot-rabbitmq-producer

6.1.2、添加依赖

修改pom.xml文件内容为如下:



    4.0.0
    com.flaw
    springboot-rabbitmq-producer
    1.0-SNAPSHOT

    
        org.springframework.boot
        spring-boot-starter-parent
        2.3.12.RELEASE
    

    
        
            org.springframework.boot
            spring-boot-starter-amqp
        
        
            org.springframework.boot
            spring-boot-starter-test
        
    


6.1.3、启动类
package com.flaw.rabbitmq;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class ProducerApplication {
    public static void main(String[] args) {
        SpringApplication.run(ProducerApplication.class);
    }
}

6.1.4、配置RabbitMQ
  1. 配置文件
    创建application.yml,内容如下:
spring:
  rabbitmq:
    host: 192.168.56.1
    port: 5672
    virtual-host: /xzk
    username: flaw
    password: flaw
  1. 绑定交换机和队列
    创建RabbitMQ队列与交换机绑定的配置类com.flaw.rabbitmq.config.RabbitMQConfig
package com.flaw.rabbitmq.config;

import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitmqConfig {
    
    public static final String ITEM_TOPIC_EXCHANGE = "springboot_item_topic_exchange";
    
    public static final String ITEM_QUEUE = "springboot_item_queue";

    
    @Bean("itemTopicExchange")
    public Exchange topicExchange(){
        return ExchangeBuilder.topicExchange(ITEM_TOPIC_EXCHANGE).durable(true).build();
    }

    
    @Bean("itemQueue")
    public Queue itemQueue(){
        return QueueBuilder.durable(ITEM_QUEUE).build();
    }

    
    @Bean
    public Binding itemQueueExchange(@Qualifier("itemQueue") Queue queue,
                                     @Qualifier("itemTopicExchange") Exchange exchange){
        return BindingBuilder.bind(queue).to(exchange).with("item.#").noargs();
    }

}

6.1.5、测试

在生产者工程springboot-rabbitmq-producer中创建测试类,发送消息:

package com.flaw.rabbitmq;

import com.flaw.rabbitmq.config.RabbitmqConfig;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.junit4.SpringRunner;

@RunWith(SpringRunner.class)
@SpringBootTest
public class RabbitMQTest {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Test
    public void test(){
        rabbitTemplate.convertAndSend(RabbitmqConfig.ITEM_TOPIC_EXCHANGE, "item.insert", "商品新增,routingKey 为item.insert");
        rabbitTemplate.convertAndSend(RabbitmqConfig.ITEM_TOPIC_EXCHANGE, "item.update", "商品修改,routingKey 为item.update");
        rabbitTemplate.convertAndSend(RabbitmqConfig.ITEM_TOPIC_EXCHANGE, "item.delete", "商品删除,routingKey 为item.delete");
    }
}

先运行上述测试程序(交换机和队列才能先被声明和绑定),然后启动消费者;在消费者工程springboot-rabbitmq-consumer中控制台查看是否接收到对应消息。
另外,也可以在RabbitMQ的管理控制台中查看到交换机与队列的绑定:

6.2、搭建消费者工程 6.2.1、创建工程

创建消费者工程springboot-rabbitmq-consumer

6.2.2、添加依赖

修改pom.xml文件内容为如下:



    4.0.0

    com.flaw
    springboot-rabbitmq-consumer
    1.0-SNAPSHOT

    
        org.springframework.boot
        spring-boot-starter-parent
        2.3.12.RELEASE
    

    
        
            org.springframework.boot
            spring-boot-starter-amqp
        
        
            org.springframework.boot
            spring-boot-starter-test
        
    


6.2.3、启动类
package com.flaw.rabbitmq;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class ConsumerApplication {
    public static void main(String[] args) {
        SpringApplication.run(ConsumerApplication.class);
    }
}

6.2.4、配置RabbitMQ

创建application.yml,内容如下:

spring:
  rabbitmq:
    host: 192.168.56.1
    port: 5672
    virtual-host: /xzk
    username: flaw
    password: flaw
6.2.5、消息监听处理类

编写消息监听器com.flaw.rabbitmq.listener.MyListener

package com.flaw.rabbitmq.listener;

import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
public class MyListener {

    
    @RabbitListener(queues = "springboot_item_queue")
    public void myListener1(String message){
        System.out.println("消费者收到的消息:"+message);
    }
}

6.2.6、测试
package com.flaw.rabbitmq;

import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.junit4.SpringRunner;


@RunWith(SpringRunner.class)
@SpringBootTest
public class ConsumerTest {

    @Test
    public void test(){
        while (true){
            // 死循环让consumer保持连接状态,一直监听队列里面的消息
        }
    }
}

7、高级特性 7.1、消息的可靠投递

在使用 RabbitMQ 的时候,作为消息发送方希望杜绝任何消息丢失或者投递失败场景。RabbitMQ 为我们提供了两种方式用来控制消息的投递可靠性模式。

  • confirm 确认模式
  • return 退回模式

rabbitmq 整个消息投递的路径为:
producer—>rabbitmq broker—>exchange—>queue—>consumer

  • 消息从 producer 到 exchange 则会返回一个 confirmCallback 。
  • 消息从 exchange–>queue 投递失败则会返回一个 returnCallback 。

我们将利用这两个 callback 控制消息的可靠性投递

创建producer项目

添加依赖



    4.0.0

    com.flaw
    spring-rabbitmq-senior-producer
    1.0-SNAPSHOT

    
        
            org.springframework
            spring-context
            5.1.7.RELEASE
        

        
            org.springframework.amqp
            spring-rabbit
            2.1.8.RELEASE
        

        
            junit
            junit
            4.12
        

        
            org.springframework
            spring-test
            5.1.7.RELEASE
        
    


创建rabbitmq.properties配置文件

rabbitmq.host=127.0.0.1
rabbitmq.port=5672
rabbitmq.username=guest
rabbitmq.password=guest
rabbitmq.virtual-host=/

创建spring-rabbitmq-producer.xml配置文件



    
    
    
    
    
    
    
    
    
    
    
    
        
            
        
    

7.1.1、确认模式

消息从 producer 到 exchange 则会返回一个 /confirm/iCallback

package com.flaw.rabbitmq;

import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;

@RunWith(SpringJUnit4ClassRunner.class)
@ContextConfiguration(locations = "classpath:spring-rabbitmq-producer.xml")
public class ProducerTest {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    
    @Test
    public void test/confirm/i() {
        // 定义回调函数
        rabbitTemplate.set/confirm/iCallback(new RabbitTemplate./confirm/iCallback() {
            
            @Override
            public void /confirm/i(CorrelationData correlationData, boolean ack, String cause) {
                System.out.println("/confirm/i方法被执行了。。。");
                if (ack) {
                    System.out.println("接收消息成功");
                }else {
                    System.out.println("接收消息失败,失败原因:"+cause);
                }
            }
        });

        // 发送消息
        //rabbitTemplate.convertAndSend("test_exchange_/confirm/i", "/confirm/i", "the message a confirm message....");
        // 失败消息发送
        rabbitTemplate.convertAndSend("test_exchange_/confirm/i111", "/confirm/i", "the message a confirm message....");
    }
}

7.2.2、退回模式

消息从 exchange–>queue 投递失败则会返回一个 returnCallback

	
    @Test
    public  void testReturn(){
        // 设置交换机处理失败消息的模式
        rabbitTemplate.setMandatory(true);

        // 定义回调函数
        rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
            
            @Override
            public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
                System.out.println("returnedMessage方法执行了。。。");

                System.out.println(message);
                System.out.println(replyCode);
                System.out.println(replyText);
                System.out.println(exchange);
                System.out.println(routingKey);
            }
        });

        // 发送信息
        rabbitTemplate.convertAndSend("test_exchange_/confirm/i", "/confirm/i111", "the message a returnCallback message....");
    }
7.2、Consumer Ack

ack指Acknowledge,确认。表示消费端收到消息后的确认方式。
有三种确认方式:

  • 自动确认:acknowledge=“none”
  • 手动确认:acknowledge=“manual”
  • 根据异常情况确认:acknowledge=“auto”,(这种方式使用麻烦,感兴趣的朋友可以研究一下)

其中自动确认是指,当消息一旦被Consumer接收到,则自动确认收到,并将相应 message 从RabbitMQ 的消息缓存中移除。但是在实际业务处理中,很可能消息接收到,业务处理出现异常,那么该消息就会丢失。如果设置了手动确认方式,则需要在业务处理成功后,调用channel.basicAck(),手动签收,如果出现异常,则调用channel.basicNack()方法,让其自动重新发送消息。

7.2.1、创建consumer项目

7.2.2、添加依赖


    4.0.0

    com.flaw
    spring-rabbitmq-senior-consumer
    1.0-SNAPSHOT

    
        
            org.springframework
            spring-context
            5.1.7.RELEASE
        

        
            org.springframework.amqp
            spring-rabbit
            2.1.8.RELEASE
        

        
            junit
            junit
            4.12
        

        
            org.springframework
            spring-test
            5.1.7.RELEASE
        
    
    

7.2.3、创建配置文件

创建 spring-rabbitmq-senior-consumersrcmainresourcesrabbitmq.properties 配置文件

rabbitmq.host=127.0.0.1
rabbitmq.port=5672
rabbitmq.username=guest
rabbitmq.password=guest
rabbitmq.virtual-host=/

创建 spring-rabbitmq-senior-consumersrcmainresourcesspring-rabbitmq-consumer.xml 配置文件



    
    

    
    

    
    

    
        
    


7.2.4、创建监听器

创建 spring-rabbitmq-senior-consumersrcmainjavacomflawrabbitmqlistenerAckListener.java 监听器

package com.flaw.rabbitmq.listener;

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
import org.springframework.stereotype.Component;


@Component
public class AckListener implements ChannelAwareMessageListener {

    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        // 获取当前消息的标签
        long deliveryTag = message.getMessageProperties().getDeliveryTag();

        try {
            System.out.println(new String(message.getBody()));
            System.out.println("处理业务逻辑");

            //模拟业务处理异常
            int i = 3/0;

            // 手动签收
            channel.basicAck(deliveryTag, true);
        } catch (Exception e) {
            //e.printStackTrace();

            
            channel.basicNack(deliveryTag,true,true);
            //channel.basicReject(deliveryTag,true);
        }
    }
}

7.2.5、测试

创建测试类 spring-rabbitmq-senior-consumersrctestjavacomflawrabbitmqConsumerTest.java

package com.flaw.rabbitmq;

import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;

@RunWith(SpringJUnit4ClassRunner.class)
@ContextConfiguration(locations = "classpath:spring-rabbitmq-consumer.xml")
public class ConsumerTest {

    @Test
    public void test(){
        while (true){
            // 死循环让consumer保持连接状态,一直监听队列里面的消息
        }
    }

}

当业务出现异常时拒绝签收当前消息,让消息重回队列进行重复执行,等待回复正常后签收。(场景:在网络波动时出现异常就拒收当前消息,让其重回队列重复获取当前消息,等网络回复正常后在签收当前消息)

7.3、消费端限流

当消费端有大量请求时(例如秒杀活动),大量的请求写入数据库时会造成数据库宕机。这时可以使用MQ存储消费端请求,让系统从MQ中拉取自身每秒最大能处理的请求量写入数据库,这样就能保证了数据库不会宕机

7.3.1、创建并配置监听
package com.flaw.rabbitmq.listener;

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
import org.springframework.stereotype.Component;


@Component
public class QosListener implements ChannelAwareMessageListener {

    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        // 休眠一秒更能直观感受
        Thread.sleep(1000);

        // 获取消息
        System.out.println(new String(message.getBody()));

        // 处理业务逻辑

        // 签收
        channel.basicAck(message.getMessageProperties().getDeliveryTag(),true);
    }
}

spring-rabbitmq-consumer.xml 配置


        

7.4、TTL

Time To Live,消息过期时间设置

管控台中设置队列TTL

代码实现(在producer项目中)
在写代码实现时先删除在管控台手动创建的队列,这样好验证代码是否创建成功

配置文件(spring-rabbitmq-producer.xml)

	
    
        
            
        
    
    
        
            
        
    

测试代码

	
    @Test
    public void testTtl(){
        
        MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
            @Override
            public Message postProcessMessage(Message message) {
                // 设置消息过期时间5秒
                message.getMessageProperties().setExpiration("5000");
                // 返回消息
                return message;
            }
        };

        // 消息单独过期
        rabbitTemplate.convertAndSend("test_exchange_ttl", "ttl.aaa", "message ttl bbb....", messagePostProcessor);

        // 注意在队列中设置过期消息不生效,消息过期后,只有消息在队列顶端,才会判断其是否过期(移除掉)。可以运行下面的代码进行验证
        
    }
7.5、死信队列

死信队列,英文缩写:DLX 。Dead Letter Exchange(死信交换机),当消息成为Dead message后,可以被重新发送到另一个交换机,这个交换机就是DLX。


消息成为死信的三种情况:

  1. 队列消息长度到达限制;
  2. 消费者拒接消费消息,basicNack/basicReject,并且不把消息重新放入原目标队列,requeue=false;
  3. 原队列存在消息过期设置,消息到达超时时间未被消费;

队列绑定死信交换机:
给队列设置参数: x-dead-letter-exchange 和 x-dead-letter-routing-key

代码实现:

spring-rabbitmq-producer.xml 配置

	
    
        
        
            
            
            
            
            
            
        
    
    
        
            
        
    
    
    
    
        
            
        
    

测试代码:

生产者测试( ProducerTest.java 测试类)

	
    @Test
    public void testDlx(){
        //  测试过期时间,死信消息
        
        //  测试超过长度限制后,消息死信
        
        // 测试消息拒收(配合消费端ConsumerTest.java测试类测试)
        rabbitTemplate.convertAndSend("test_exchange_dlx","test.dlx.haha","我是一条nack消息,我会死吗?");
    }

消费者监听器 DlxListener.java

package com.flaw.rabbitmq.listener;

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
import org.springframework.stereotype.Component;


@Component
public class DlxListener implements ChannelAwareMessageListener {

    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();

        try {
            // 接收转换消息
            System.out.println(new String(message.getBody()));
            // 处理业务逻辑
            System.out.println("处理业务逻辑...");
            // 模拟业务出现异常
            int i = 3/0;
            // 手动签收
            channel.basicAck(deliveryTag,true);
        } catch (Exception e) {
            //e.printStackTrace();
            System.out.println("出现异常,拒绝接受");
            // 拒绝签收,不重回队列 requeue=false
            channel.basicNack(deliveryTag,true,false);
        }
    }
}

spring-rabbitmq-consumer.xml 配置

	
        
    
7.6、延迟队列

延迟队列,即消息进入队列后不会立即被消费,只有到达指定时间后,才会被消费。

需求:

  1. 下单后,30分钟未支付,取消订单,回滚库存。
  2. 新用户注册成功7天后,发送短信问候。

实现方式:

  1. 定时器
  2. 延迟队列


很可惜,在RabbitMQ中并未提供延迟队列功能。但是可以使用:TTL+死信队列 组合实现延迟队列的效果。

代码实现:

生产者测试( ProducerTest.java 测试类)

	@Test
    public void testDelay() throws Exception {
        SimpleDateFormat sdf=new SimpleDateFormat("yyyy年MM月dd日HH:mm:ss SSS");
        // 发送订单消息。 模拟在订单系统中,下单成功后,10秒内未支付该订单就会进入死信队列
        rabbitTemplate.convertAndSend("order_exchange","order.msg","订单信息: id=1,time="+sdf.format(new Date()));

        // 打印10秒倒计时 (10秒后订单未处理就会自动发送到死信队列,当启动Consumer测试时,发送的订单信息会被消费不会进入死信队列)
        for (int i = 10; i > 0; i--) {
            System.out.println(i+"...");
            Thread.sleep(1000);
        }
    }

消费者监听器(OrderListener.java)

package com.flaw.rabbitmq.listener;

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
import org.springframework.stereotype.Component;

@Component
public class OrderListener implements ChannelAwareMessageListener {

    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();

        try {
            // 接收转换消息
            System.out.println(new String(message.getBody()));
            // 处理业务逻辑
            System.out.println("处理业务逻辑...");
            System.out.println("订单支付完成...");
            // 手动签收
            channel.basicAck(deliveryTag,true);
        } catch (Exception e) {
            //e.printStackTrace();
            System.out.println("出现异常,拒绝接受");
            // 拒绝签收,不重回队列 requeue=false
            channel.basicNack(deliveryTag,true,false);
        }
    }
}

spring-rabbitmq-consumer.xml 配置

	
        
    

当消费者没启动时,生产者测试产生的订单会在10秒后进入死信队列
启动消费者模拟完成支付业务,生产者产生的订单就会被消费

转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/354494.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

版权所有 (c)2021-2022 MSHXW.COM

ICP备案号:晋ICP备2021003244-6号