RabbitMQ-消息队列(小白入门篇)
1. 部署
这里我们采用docker部署
docker run -d \
-e RABBITMQ_DEFAULT_USER=YLZ \
-e RABBITMQ_DEFAULT_PASS=123 \
-v mq-plugins:/plugins \
--name mq \
-p 15672:15672 \
-p 5672:5672
rabbitmq:4-management
这里先剧透一下,后面创建mq里的资源一般是通过代码,所以这里没有挂载用于持久化的数据卷
访问15672端口,这是mq的控制台

2. 资源类型
这是RabbitMQ的架构图,主要聚焦在exchange、queue、Binding、VirtualHost这三种

1. Exchange-交换机
点击Exchange选项卡,可以看到默认的交换机

生产者只需要将消息发送到交换机就行,然后交换机会把消息转发给与其绑定的队列,不过交换机有几种类型,其转发的规则也有所不同,我们主要介绍三种:fanout、direct、topic
1. fanout
这是广播交换机,会把收到的消息分发给所有与其绑定的队列
2. direct
对比fanout,direct可以做到将消息转发到指定的队列,只需要在绑定队列的时候指定队列的routingKey,在发消息的时候携带上目标队列的routingKey,就可以实现向指定队列发消息;如果没携带routingKey的话,则会把消息丢掉;并且,一个队列可以可以绑定多个routingKey,两个队列可以绑定同一个routingKey
3. topic
对比direct,topic在匹配routingKey的时候,routingKey的格式是用“ . ”隔开,比如aaa.bbb.ccc,在topic眼里就是aaa中的bbb中的ccc,并且可以有通配符,其中:
- “ # ” 是0个或多个字母
- “ . ” 是一个单词
比如aaa.#,可以是aaa.bbb或aaa.ccc或aaa.ccc.bbb;也可以写#.bbb
2. Queue-队列
点击Queues and Streams选项卡,默认是没有队列的,我们手动创建一个:选择Add a new queue,然后填队列名,选择type(一般是Classic),然后点Add queue按钮

WorkQueue:队列中的每个消息都只能被消费一次,被消费者取出来后就没有了
3. Binding-绑定
来到交换机页面,点击要操作的交换机

然后选择binding,这里用的是fanout,所以可以不用填Routing key,如果是别的交换机就需要!

4. VirtualHost-虚拟主机
点击Admin选项卡,点击右侧菜单的Virtual Hosts;可以看到默认的虚拟主机是“/”,可以在Add a new virtual host新增虚拟主机

这里新增一个“aaa”虚拟主机
可以点进去配置该虚拟主机可以被其他用户使用

我们回到Exchange页,可以发交换机多了,这是新增的虚拟主机创建的默认交换机

如果只聚焦于一个虚拟主机的话,我们可以选择右上角的Virtual host

3. java客户端
程序操作rabbitMQ需要通过amqp协议,这有专门的amqp的SDK可以用,不过spring-boot有对其整合,所以我们直接用spring-boot管理的amqp依赖
1. 引入依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
2. 配置application.yml
spring:
rabbitmq:
# rabbitmq的地址
host: 192.168.233.128
# rabbitmq的端口号
port: 5672
# 虚拟主机
virtual-host: aaa
# 用户名
username: YZL
# 密码
password: 123
3. 发消息
springboot有默认的rabbitMQ客户端实现——RabbitTemplate,并且自动装配了,可以直接注入,就像RedisTemplate一样。以下是四种生产者发消息的示例
@SpringBootTest
class JavaDemoApplicationTests {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 直接向队列发送消息
*/
@Test
void simplePubQueue() {
String queueName = "simple.queue";
String message = "hello fuck";
rabbitTemplate.convertAndSend(queueName,message);
}
/**
* 向 fanout交换机 发送消息
* fanout交换机 是广播, 将消息发送给所有绑定的队列
*/
@Test
void fanoutExchange() {
String exchange = "lzy.fanout";
String message = "hello fuck";
rabbitTemplate.convertAndSend(exchange,message);
}
/**
* 向 direct交换机 发送消息
* direct交换机 是路由(toutingKey), 根据路由键将消息发送给指定的队列
**/
@Test
void directExchange() {
String exchange = "lzy.direct";
String message = "hello direct-a";
rabbitTemplate.convertAndSend(exchange,"a",message);
}
/**
* 向 topic交换机 发送消息
* topic交换机 是路由(toutingKey), 根据路由键将消息发送给指定的队列
* 与 direct交换机 不同的是, topic 的路由可以使用通配符
* 路由键可以使用多个单词, 用 . 分隔, 用 * 匹配一个单词, 用 # 匹配 0个或多个单词
**/
@Test
void topicExchange() {
String exchange = "lzy.topic";
String message = "hello all news";
String aCountry = "hello a";
// 向所有国家的news队列发送消息
rabbitTemplate.convertAndSend(exchange,"#.news",message);
// 向a国家所有队列发送消息
rabbitTemplate.convertAndSend(exchange,"a.#",aCountry);
}
}
4. 消费消息
创建监听消息类,需要注册为Bean,这里直接用@Component注解,然后写消费方法,只需要用@RabbitListener(queues = "queue.name")注解,就能监听队列的消息
@Component
public class MyListener {
@RabbitListener(queues = "simple.queue")
public void simpListenQueue(String message){
System.out.println("监听simple.queue: "+message);
}
}
可以发现:发消息只需要关注发到哪个交换机,收消息只需要关注哪个队列,从而实现了业务解耦
5. 创建资源
前面说到可以依靠程序来创建mq的资源
1. 通过注解声明
// 声明式创建交换机、队列、绑定关系, 并监听
@RabbitListener(bindings = @QueueBinding(
exchange = @Exchange(value = "ylz.topic", type = ExchangeTypes.TOPIC),
value = @Queue(value = "ylz.queue"),
key = "ylz.topic.key"))
public void fanoutListenQueue(String message){
System.out.println("ylz.queue: "+message);
}
这里使用@RabbitListener注解,里面的属性需要声明:绑定关系、交换机、队列,创建后就开始监听队列消息。如果是fanout交换机的话就不用指定key,type的默认值是direct交换机
2. 通过类声明
spring-amqp提供了几个类,用于创建交换机、队列、绑定关系的bean;这里直接展示示例代码,不知道要写啥,主要在配置中声明
@Configuration
public class Config {
@Bean
public MessageConverter messageConverter(){
return new Jackson2JsonMessageConverter();
}
@Bean
public FanoutExchange fanoutExchange(){
// 通过构造器创建
// return new FanoutExchange("ylz.fanout");
// 通过 建造者 创建
return ExchangeBuilder
.fanoutExchange("ylz.fanout")
.build();
}
@Bean
public Queue queue1(){
// 通过构造器创建
// return new Queue("ylz.queue1");
// 通过 建造者 创建
return QueueBuilder
.durable("ylz.queue1")
.build();
}
@Bean
public Binding binding(Queue queue1, FanoutExchange fanoutExchange){
// 绑定队列到交换机, 如果不是fanout交换机的话, 后面还可以.with(routingKey)给路由
return BindingBuilder.bind(queue1).to(fanoutExchange);
}
}
如果不是fanout的话,绑定时还需要注解key(routingKey),如下
@Bean
public Binding binding(Queue queue1, DirectExchange directExchange){
// 绑定队列到交换机, 如果不是fanout交换机的话, 后面还可以.with(routingKey)给路由
return BindingBuilder.bind(queue1).to(fanoutExchange).with("a.news");
}
现在启动,查看mq控制台,查看交换机与绑定信息


6. 消息格式
发送消息是可以发送java对象的,他会序列化成二进制发送,消费者接受对象后反序列化
// 对象
@Data
@AllArgsConstructor
@NoArgsConstructor
public class User implements Serializable {
private static final long serialVersionUID = 1L;
// 负责接受前端传来的数据和把数据响应给前端
private Integer id; // 用包装类以便确认数据是否有接到
private String name;
private String password;
private Integer sex;
private Integer age;
private String hobby;
}
// 生产者
public void simplePubQueue() {
String queueName = "ylz.queue";
User message = new User(1,"ylz","123456",1,18,"football");
rabbitTemplate.convertAndSend(queueName,message);
}
// 消费者
@RabbitListener(bindings = @QueueBinding(
exchange = @Exchange(value = "ylz.topic", type = ExchangeTypes.TOPIC),
value = @Queue(value = "ylz.queue"),
key = "ylz.topic.key"))
// 注意参数类型,要与消息类型一样才能反序列化
public void fanoutListenQueue(User message){
System.out.println("ylz.queue: "+message.toString());
}
发送出的消息,可以在控制台的队列页面的Get messages看到

可以看到消息内容都看不懂,这就是序列化后的java对象,占用比较大,我们接下来要把java对象序列化成json格式来发消息,这样占用更小
1. 引入json序列化依赖
生产者和消费者服务都要引入
<dependency>
<groupId>com.fasterxml.jackson.dataformat</groupId>
<artifactId>jackson-dataformat-xml</artifactId>
</dependency>
2. 实现MessageConverter接口
rabbitMQ的消息转换都是通过这个接口来执行,他有其他的实现类,其中就包括json格式的实现类,默认是java对象序列化
我们要找的json格式的实现类如图

将该接口注册为bean
@Configuration
public class MqConfig {
@Bean
public MessageConverter messageConverter(){
return new Jackson2JsonMessageConverter();
}
}
然后重启再发消息,回到控制台,可以看到消息变成json格式,且占用大大降低

至此,你已经能使用消息队列的基本操作了
更多推荐


所有评论(0)