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的架构图,主要聚焦在exchangequeueBindingVirtualHost这三种

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交换机的话就不用指定keytype的默认值是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格式,且占用大大降低

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

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐