栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 云计算 > 云平台

RabbitMQ

云平台 更新时间: 发布时间: IT归档 最新发布 模块sitemap 名妆网 法律咨询 聚返吧 英语巴士网 伯小乐 网商动力

RabbitMQ

RabbitMQ

文章目录
    • RabbitMQ
      • 1、MQ的基本概念
      • 2、RabbitMQ
      • 3、RebbitMQ的安装
      • 4、Rabbitmq的工作模式
        • 4.1 简单模式
        • 4.2 工作队列模式
        • 4.3 发布订阅模式
        • 4.4 路由模式
        • 4.5 主题模式
      • 5、Spring整合RabbitMQ
      • 6、Spring Boot整合RabbitMQ
    • RabbitMQ 高级特性
      • 1、消息的可靠投递
      • 2、Consumer Ack
      • 3、消费端限流
      • 4、TTL 过期时间
      • 5、死信队列
      • 6、延迟队列
      • 7、日志与监控
    • RabbitMQ 应用问题
      • 1、消息可靠性保障
      • 2、消息幂等性保障

1、MQ的基本概念

概念

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

图示:

  • MQ,消息队列,存储信息的中间件
  • 分布式系统通信的两种方式:直接远程调用和借助MQ消息中间件完成通信
  • 发送方称为生产者,接收方称为消费者

MQ的优缺点

  • 优点:
    • 解耦
    • 异步提速,提高效率
    • 削峰填谷,提高系统的稳定性
  • 缺点:
    • 系统的可用性降低 【整合的技术越多,维护的就越多】
    • 系统复杂程度提高
    • 一致性问题

MQ的适用场景

  • 生产者不需要从消费者处获取反馈,接口的返回值为空值,生产者只需发送消息到消息队列,不需关注后面消息是如何消费的
  • 容许短暂的不一致性
  • 确实要用到MQ的功能

常见的四大MQ

RabbitMQActiveMQRocketMQKafka
社区RabbitApache阿里巴巴Apache
开发语言ErlangjavajavaScala&java
协议支持AMQP、SMTP、XMPP、STOMPOpenWire、STOMP、AMQP自定义自定义
客户端支持语言java、Erlangjava、c、c++java、c++java
单机吞吐量万级(其次)万级(最差)十万级(最好)十万级(次之)
消息延迟微秒级毫秒级毫秒级毫秒级
2、RabbitMQ

AMQP

AMQP是一个网络协议,是应用层协议的一个开放标准,为面向消息中间件设计,是一个消息中间件开发的一种标准

rabbitMQ简介

rabbitMQ是基于AMQP标准开发的一个消息中间件

  • 一个Broker中可以与多个虚拟机
  • 一个虚拟机由多个交换机和多个队列构成
  • 交换机与队列之间有绑定关系
  • RabbitMQ有6种工作模式:简单模式 、工作队列模式 、发布订阅模式 、路由模式 、主题模式 、rpc模式
  • 官网:https://www.rabbitmq.com/

JMS

  • JMS即java消息服务应用程序接口,是一个java平台中关于面向消息中间件的API
  • 是javaEE规范的一种,类上JDBC
  • 许多消息中间件都实现了JMS规范
3、RebbitMQ的安装

linux上安装rebbitMQ

  • 下载 erlang-23.2.7-2.el7.x86_64.rpm 包,环境【注意版本要与rebbitmq版本匹配】
  • 下载 rabbitmq-server-3.9.16-1.el8.noarch.rpm包 【注意版本要与erlang版本匹配】
  • 默认安装在 usr/lib/rabbitmq/rabbitmq_server-3.9.16/sbin/
  • 配置文件在 usr/lib/rabbitmq/lib/
  • 开启可视化插件:rabbitmq-plugins enable rabbitmq_management
  • 开启端口:
    • systemctl start firewalld 开启防火墙
    • firewall-cmd --add-port=15672/tcp --permanent 开启端口
    • firewall-cmd --list-ports 查看
  • sbin目录下启动 rabbitmq-server
  • 访问:虚拟机地址:15672/

由于guest guest 用户只能在本机访问,异机登录需要自己设置一个用户: 【启动rabbitmq服务后,执行以下命令】

  • **rabbitmqctl add_user admin 123 ** 添加用户和密码
  • rabbitmqctl set_user_tags admin administrator 将添加的用户设置为超级管理员
  • rabbitmqctl set_permissions -p / admin “." ".” “.*” 设置权限,这样是guest用户所拥有的全部权限admin都有

  • 这样就成功了,可用设置的用户登录,可以在它的界面进行消息的发送与消费,也可用java代码向消息队列发送和消费消息
4、Rabbitmq的工作模式

Rabbitmq有6中工作模式:简单模式 、工作队列模式 、发布订阅模式 、路由模式 、主题模式 、rpc模式

4.1 简单模式

简单模式是指一个消息的发送方将消息发送到消息中间件,默认的交换机,一个消费者监听着该消息队列消费其中的消息

生产者:

向消息中间件发送消息的一方

  1. 创建就连接工厂
  2. 设置参数
  3. 创建连接 【确保15672和5672的端口是否已经打开】
  4. 创建channel
  5. 创建队列
  6. 发送消息
  7. 关闭资源

依赖:


    
    
        com.rabbitmq
        amqp-client
        5.14.1
    


    8
    8


代码

package com.qiumin;

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

import java.io.IOException;
import java.util.concurrent.TimeoutException;


public class Producer_HW {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory.setHost("自己的主机");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //5.创建队列
//queueDeclare(String queue, boolean durable, boolean exclusive, boolean autoDelete, Map arguments)
        //队列名称, 是否持久化 , 是否内部使用  ,是否自动删除 , 参数
        channel.queueDeclare("hello_word",true,false,false,null);
        //6.发送消息
        String body = "欢迎使用消息中间件 成功";
  // public void basicPublish(String exchange, String routingKey, BasicProperties props, byte[] body)
        //交换机名称该模式下默认为“” ,路由key该模式下为队列名称 , 参数 , 内容需要是字节型
        channel.basicPublish("","hello_word",null,body.getBytes());
        //7.关闭资源
        channel.close();
       connection.close();
    }
}

访问: ip+端口

消费者

从消息队列中拿取消息进行消费

依赖:


    
    
        com.rabbitmq
        amqp-client
        5.14.1
    


    8
    8


代码

package com.qiumin;

import com.rabbitmq.client.*;

import java.io.IOException;
import java.util.concurrent.TimeoutException;


public class Consumer_HW {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory.setHost("自己的主机");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //4.创建队列, 该方法如果所建的队列存在就不会再创建
        channel.queueDeclare("hello_word",true,false,false,null);
        //5.消费
        //回调自动确认方法
        DefaultConsumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.out.println(new String(body));
            }
        };
        //第二参数为是否自动确认
        channel.basicConsume("hello_word",true,consumer);
        //7.注意消费者一般不需要关闭资源,一直监听队列,有消息就进行消费
       
    }
}

访问: ip+端口

4.2 工作队列模式

一个消费者向队列中发送消息,多个消费者监听同一个队列,进行消息的消费

  • 建立两个消费者,与上面的消费者一模一样,先启动两个消费者进行监听同一个队列,生产端再向该队列中发送消息
  • 可以看到两个消费者轮询的方式消费队列中的消息
4.3 发布订阅模式

发布订阅模式

发布订阅模式中多了一个Exchange交换机的角色,而交换机有常见的三种类型:

  1. Fanout: 广播,将消息交给所有绑定到交换机的队列
  2. Direct: 定向,将消息发送给符合指定 routing key 的队列
  3. Topic: 通配符,将消息发给符合通配符 routing pattern 的队列 【可以兼容以上两种模式,* 代表一个单词,# 代表0个或多个单词】

Exchange只负责转发消息,不具备存储消息的能力,因此如果没有任何队列与交换机绑定,消息就会丢失

生产者

步骤:

  • 创建连接工厂
  • 设置连接相关参数
  • 创建连接
  • 创建channel
  • 创建交换机
  • 创建队列
  • 交换机与队列绑定
  • 发送消息
  • 关闭资源

代码

package com.qiumin;

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

import java.io.IOException;
import java.util.concurrent.TimeoutException;


@SuppressWarnings("all")
public class Producer_pub {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory3 = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory3.setHost("自己的主机");
        connectionFactory3.setPort(5672);
        connectionFactory3.setVirtualHost("/");
        connectionFactory3.setUsername("admin");
        connectionFactory3.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory3.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //5.创建交换机
        //exchangeDeclare(String exchange,  交换机名字
        // BuiltinExchangeType type, 类型
        // boolean durable,  是否持久化
        // boolean autoDelete, 是否自动删除
        // boolean internal,  是否内部使用
        // Map arguments) 参数
        channel.exchangeDeclare("Test_fanout", BuiltinExchangeType.FANOUT,true,false,false,null);
        //6.创建队列
        channel.queueDeclare("queue1",true,false,false,null);
        channel.queueDeclare("queue2",true,false,false,null);
        //7.交换机与队列的绑定
        //参数 queue:队列名字  exchange:交换机名字  routingKey:路由键  FANOUT类型的交换机不需要设置路由,默认路由""
        channel.queueBind("queue1","Test_fanout","");
        channel.queueBind("queue2","Test_fanout","");
        //8.发送消息
        for (int i = 0; i < 10; i++) {
            String body = "消息"+i;
            channel.basicPublish("Test_fanout","",null,body.getBytes());
        }
        //9.关闭资源
        channel.close();
        connection.close();
    }
}

消费者

package com.qiumin;

import com.rabbitmq.client.*;

import java.io.IOException;
import java.util.concurrent.TimeoutException;


public class Consumer_sub {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory.setHost("自己的主机");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //4.创建队列, 该方法如果所建的队列存在就不会再创建
        channel.queueDeclare("queue1",true,false,false,null);
        //5.消费消息
        DefaultConsumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.out.println(new String(body)+"queue1");
            }
        };
        channel.basicConsume("queue1",true,consumer);  //监听队列1
        // channel.basicConsume("queue2",true,consumer);  另一个消费者这里监听队列2
    }
}
  • 两个队列都有交换机发来的相同数量相同内容的消息
  • 这就是fanout 广播类型的交换机
4.4 路由模式

路由模式,交换机根据routkey来转发消息到指定的队列

说明

  • 队列和交换机的绑定不能任意了,而是要指定一个routingKey 路由key
  • 消息发送到交换机是必须指定一个routingKey
  • 交换机不会把消息转发到每个队列,而是与routingKey匹配的队列
  • 交换机类型为 BuiltinExchangeType.DIRECT

生产者

  • 按routingKey将交换机与队列进行绑定
package com.qiumin;

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

import java.io.IOException;
import java.util.concurrent.TimeoutException;


public class Producer_RK {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory3 = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory3.setHost("自己的主机");
        connectionFactory3.setPort(5672);
        connectionFactory3.setVirtualHost("/");
        connectionFactory3.setUsername("admin");
        connectionFactory3.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory3.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //5.创建交换机
        channel.exchangeDeclare("Test_direct", BuiltinExchangeType.DIRECT,true,false,false,null);
        //6.创建队列
        channel.queueDeclare("queue_RK1",true,false,false,null);
        channel.queueDeclare("queue_RK2",true,false,false,null);
        //7.交换机与队列的绑定
        //参数 queue:队列名字  exchange:交换机名字  routingKey:路由键  FANOUT类型的交换机不需要设置路由,默认路由""
        channel.queueBind("queue_RK1","Test_direct","info"); //路由key为info,带有info的转发到该队列
        channel.queueBind("queue_RK2","Test_direct","warning");//路由key为warning,带有info的转发到该队列
        //8.发送消息
        String body1 = "info....信息";
        String body2 = "warning....信息";
        //发送路由key为info的消息会转发到匹配的的队列queue_RK1中
        channel.basicPublish("Test_direct","info",null,body1.getBytes());
        channel.basicPublish("Test_direct","info",null,body1.getBytes());
        channel.basicPublish("Test_direct","info",null,body1.getBytes());
        //发送路由key为warning的消息会转发到匹配的的队列queue_RK2中
        channel.basicPublish("Test_direct","warning",null,body2.getBytes());
        channel.basicPublish("Test_direct","warning",null,body2.getBytes());
        channel.basicPublish("Test_direct","warning",null,body2.getBytes());
        //9.关闭资源
        channel.close();
        connection.close();
    }
}

消费者

package com.qiumin;

import com.rabbitmq.client.*;

import java.io.IOException;
import java.util.concurrent.TimeoutException;


public class Consumer_RK {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory.setHost("自己的主机");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //4.创建队列, 该方法如果所建的队列存在就不会再创建
        channel.queueDeclare("queue_RK1",true,false,false,null);
        //5.消费消息
        DefaultConsumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.out.println(new String(body)+"queue1");
            }
        };
        channel.basicConsume("queue_RK1",true,consumer);  //监听接收带有路由key info消息的队列queue_RK1
       // channel.basicConsume("queue_RK2",true,consumer);  监听接收带有路由key warning消息的队列queue_RK2
    }
}
  • 不同的路由key的消息转发到与之相匹配的队列中
4.5 主题模式

Topic主题模式

  • 交换机类型为 Topic
  • 通过统配符转发到指定队列
  • admin.# 以admin开头的 #代表0个或多个单词
  • *.list * 代表1个单词

生产者

package com.qiumin;

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

import java.io.IOException;
import java.util.concurrent.TimeoutException;


public class Producer_TP {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory3 = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory3.setHost("自己的主机");
        connectionFactory3.setPort(5672);
        connectionFactory3.setVirtualHost("/");
        connectionFactory3.setUsername("admin");
        connectionFactory3.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory3.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //5.创建交换机
        //exchangeDeclare(String exchange,  交换机名字
        // BuiltinExchangeType type, 类型
        // boolean durable,  是否持久化
        // boolean autoDelete, 是否自动删除
        // boolean internal,  是否内部使用
        // Map arguments) 参数
        channel.exchangeDeclare("Test_topic", BuiltinExchangeType.TOPIC,true,false,false,null);
        //6.创建队列
        channel.queueDeclare("queue_TP1",true,false,false,null);
        channel.queueDeclare("queue_TP2",true,false,false,null);
        //7.交换机与队列的绑定
        //参数 queue:队列名字  exchange:交换机名字  routingKey:路由键  FANOUT类型的交换机不需要设置路由,默认路由""
        channel.queueBind("queue_TP1","Test_topic","#.haha");  //通配符方式的key  haha结尾的  #0或多个单词
        channel.queueBind("queue_TP2","Test_topic","heihei.*"); //通配符方式的key heihei开头的 * 1个单词
        //8.发送消息
        String body3 = "haha结尾的....信息";
        String body4 = "heihei开头的....信息";

        channel.basicPublish("Test_topic","info.qq.haha",null,body3.getBytes());
        channel.basicPublish("Test_topic","name.qq.haha",null,body3.getBytes());
        channel.basicPublish("Test_topic","heihei.admin",null,body4.getBytes());
        //9.关闭资源
        channel.close();
        connection.close();
    }
}

消费者

package com.qiumin;

import com.rabbitmq.client.*;

import java.io.IOException;
import java.util.concurrent.TimeoutException;


public class Consumer_TP {
    public static void main(String[] args) throws IOException, TimeoutException {
        //1.创建就连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接参数参数
        connectionFactory.setHost("自己的主机");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");
        //3.创建连接
        Connection connection = connectionFactory.newConnection();
        //4.创建channel
        Channel channel = connection.createChannel();
        //4.创建队列, 该方法如果所建的队列存在就不会再创建
        channel.queueDeclare("queue_TP1",true,false,false,null);
        //5.消费消息
        DefaultConsumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.out.println(new String(body)+"queue_TP1");
            }
        };
        channel.basicConsume("queue_TP1",true,consumer); //监听queue_TP1
        // channel.basicConsume("queue_TP2",true,consumer); //另一个监听queue_TP2
    }
}
5、Spring整合RabbitMQ

使用spring整合raabbitmq

创建工程—> 添加依赖 —>配置 【关键】—>代码—>运行

  • 依赖

    
        org.springframework
        spring-context
        5.3.17
    
    
        org.springframework.amqp
        spring-rabbit
        2.4.3
    
    
        org.springframework
        spring-test
        5.3.17
    
    
        junit
        junit
        4.12
    

  • properties
rabbitmq.host=自己的主机ip
rabbitmq.port=5672
rabbitmq.username=admin
rabbitmq.password=admin
rabbitmq.virtual-host=/
  • xml



    
    

    
    
    
    
    
    

    
    
    
    
    
    
        
            
            
        
    

    
    
    
    
    
    
        
            
            
        
    

    
    


  • 测试
package com.qiumin;

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) //运行器,该注解不能丢,否则rabbitTemplate为空,报空指针异常
@ContextConfiguration(locations = "classpath:spring-rabbitmq-producer.xml")
public class Producer_Spring {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Test
    public void  Test(){
        //简单模式
        rabbitTemplate.convertAndSend("spring_queue","简单类型的交换机");
    }

    @Test
    public void  TestFanout(){
        //简单模式
        rabbitTemplate.convertAndSend("fanout_exchange","","广播类型的交换机 fanout");
        rabbitTemplate.convertAndSend("fanout_exchange","","广播类型的交换机 fanout2");
    }

    @Test
    public void  TestTopic(){
        //简单模式
        rabbitTemplate.convertAndSend("topic_exchange","admin.haha","主题类型类型的交换机 topic_exchange");
        rabbitTemplate.convertAndSend("topic_exchange","qiumin.admin.haha","主题类型的交换机 topic_exchange");
        rabbitTemplate.convertAndSend("topic_exchange","heihei.qiumin","主题类型类型的交换机 topic_exchange");
    }
}

消费者

  • 依赖

    
        org.springframework
        spring-context
        5.3.17
    
    
        org.springframework.amqp
        spring-rabbit
        2.4.3
    
    
        org.springframework
        spring-test
        5.3.17
    
    
        junit
        junit
        4.12
    

  • properties
rabbitmq.host=自己的主机ip
rabbitmq.port=5672
rabbitmq.username=admin
rabbitmq.password=admin
rabbitmq.virtual-host=/
  • 监听类
package com.qiumin.listener;

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




public class TopicLisener implements MessageListener {
    @Override
    public void onMessage(Message message) {
        System.out.println(new String(message.getBody())+"queue_topic1"); //监听队列queue_topic1
    }
}
//====================================================================================================
package com.qiumin.listener;

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




public class TopicLisener2 implements MessageListener {
    @Override
    public void onMessage(Message message) {
        System.out.println(new String(message.getBody())+"queue_topic2");//监听队列queue_topic2
    }
}

  • xml



    
    

    
    

    
    
    

    
    
        
        
    


  • 测试
@RunWith(SpringJUnit4ClassRunner.class) //运行器
@ContextConfiguration(locations = "classpath:spring-rabbitmq-consumer.xml")
public class Consumer_spring {

    @Test
    public void  Test(){
     while (true){}
    }
    
}

总结:

  1. 配置文件创建交换机队列及绑定关系 注入rabbitTamplate进行操作
  2. 消费端写监听处理类通过配置文件绑定监听的队列
6、Spring Boot整合RabbitMQ

Spring Boot整合RabbitMQ

创建项目 —>导pom —>改yaml —>配置类【关键】 —>测试

生产者

  • pom

    
        org.springframework.boot
        spring-boot-starter-web
        2.6.5
    
    
        org.springframework.boot
        spring-boot-starter-test
        2.6.5
    
    
        org.springframework.boot
        spring-boot-starter-amqp
        2.6.6
    

  • yaml
server:
  port: 8081

spring:
  rabbitmq:
    addresses: #自己的主机ip
    port: 5672
    username: admin
    password: admin
    virtual-host: /   #连接相关的参数信息
  • 配置类
package com.qiumin.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 EXCHANGE_NAME="boot_topic_exchange";
    public static final String QUEUE_NAME="boot+queue";

    //交换机exchange
    @Bean("exchange")
    public Exchange bootExchange(){
        return ExchangeBuilder.topicExchange(EXCHANGE_NAME).durable(true).build();
    }

    //队列 queue
    @Bean("queue")
    public Queue bootQueue(){
        return QueueBuilder.durable(QUEUE_NAME).build();
    }

    //交换机与队列的绑定
    @Bean
    public Binding bindingExchangeQueue(@Qualifier("exchange") Exchange exchange,@Qualifier("queue") Queue queue){
        return BindingBuilder.bind(queue).to(exchange).with("boot.#").noargs();
    }
}
  • 测试
@SpringBootTest
public class TestBootTopic {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Test
    public void test(){
        rabbitTemplate.convertAndSend("boot_topic_exchange","boot.admin","spring boot 整合 rabbitmq");
    }
}

消费者 @RabbitListener(queues = “queue的名字”)

  • pom

    
        org.springframework.boot
        spring-boot-starter-web
        2.6.5
    
    
        org.springframework.boot
        spring-boot-starter-test
        2.6.5
    
    
        org.springframework.boot
        spring-boot-starter-amqp
        2.6.6
    

  • yaml
server:
  port: 8080

spring:
  rabbitmq:
    addresses: #自己的主机ip
    port: 5672
    username: admin
    password: admin
    virtual-host: /   #连接相关的参数信息
  • 配置类
package com.qiumin;

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



@Component
public class RabbitmqLisener {

    @RabbitListener(queues = "boot+queue")
    public void LisenerQueue(Message message){
        System.out.println(new String(message.getBody()));
    }
}

  • 测试
@SpringBootTest
public class TestTopic {

    @Test
    public void test(){

    }
}

总结:

  1. 通过配置文件创建交换机队列和之间的绑定关系,注入rabbitTemplate进行操作
  2. 消费端直接使用@RabbitListener(queues = “boot+queue”) 监听指定队列
RabbitMQ 高级特性 1、消息的可靠投递

可靠性投递是指作为发送方要拒绝消息投递失败的场景,保证消息能够正确的投递到exchange

控制消息的可靠性投递的两种模式模式:

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

消息的投递路径为:

producer —>rabbitmq broker —>exchange —>queue —>consumer

  • 消息从producer到exchange会返回一个confirmCallback
  • 消息从exchange到queue投递失败会返回一个returnCallback
  • 利用这两个回调函数判断消息的可靠性投递

confirm确认模式 【spring】

  • 开启

  • 代码 rabbitTemplate.setConfirmCallback()设置回调
 

return退回模式

  • 开启

  • 设置 rabbitTemplate.setMandatory(true)及回调函数
 @Test
    public void  TestTopic(){
        //主题模式模式
        
        //出错回调
        rabbitTemplate.setMandatory(true);
        
        //设置回调函数
        rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() {
            @Override
            public void returnedMessage(ReturnedMessage returnedMessage) {
                System.out.println("returnedMessage方法被执行了----出现错误");
                System.out.println(returnedMessage);
            }
        });

        rabbitTemplate.convertAndSend("topic_exchange","admin.hahl","主题类型类型的交换机 topic_exchange");                                    //key错误,不能到达queue
        rabbitTemplate.convertAndSend("topic_exchange","qiumin.admin.hahl","主题类型的交换机 topic_exchange");                                    //key错误,不能到达queue                          
        rabbitTemplate.convertAndSend("topic_exchange","heihel.qiumin","主题类型类型的交换机 topic_exchange");									 //key错误,不能到达queue
    }

总结:

  • 开启publisher-returns=“true”
  • 使用confirm确认模式,设置对应的回调函数即可 【producer—>exchange阶段】
  • 使用return回退模式,出错时才会调用,设置开启 rabbitTemplate.setMandatory(true); ,设置回调函数【exchange—>queue阶段】
  • 可以使用事务的机制 channel的方法但性能较差,推荐使用前面两种方式
    • channel.txSelect:设置为事务模式
    • channel.txCommit:提交
    • channel.txRollback:回滚
2、Consumer Ack

Ack是指Acknowledge确认,表示消费端收到消息的确认方式

三种方式:

  • 自动确认:acknowledge=“none” 一旦收到消息就自动确认收到,不管业务是否处理成功,消息被移除
  • 手动确认:acknowledge=“manual” 收到消息后业务处理成功,手动调用channel.basicAck()确认收到,业务处理出现异常 channel.basicNack() 拒收,方法中可以设置参数来约束是否将消息返回到queue,重新发送消息
  • 根据异常情况确认:acknowledge=“auto” (处理相对麻烦,根据异常的类型,具体情况具体处理)

手动确认签收 【spring】

  • 在rabbit:listener-container 中设置 acknowledge=“manual”

  • 监听类设置
package com.qiumin.listener;

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


public class AckListener implements ChannelAwareMessageListener {
    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        //接收到消息
        System.out.println(new String(message.getBody()));
       try {
           //处理业务
          // int i=1/0;  制造异常
           int i = 1+1;
           System.out.println("处理业务..."+":"+i);
           //处理正常,手动签收
           channel.basicAck(deliveryTag,true); //true代表签收多个消息
       }catch (Exception e){
           //处理出现异常拒收
           System.out.println("异常");
           channel.basicNack(deliveryTag,true,true);
           //第一个true代表拒收多个消息,第二个true代表是否将消息重回队列,重新发送
       }

    }
}
  • 实现 ChannelAwareMessageListener 接口
  • 调用channel的方法

总结:

  • 在rabbit:listener-container 中设置 acknowledge=“manual”。
  • 没有异常调用 channel.basicAck(deliveryTag,true); //true代表签收多个消息。
  • 出现异常调用 channel.basicNack(deliveryTag,true,true); //第一个true代表拒收多个消息,第二个true代表是否将消息重回队列,重新发送。
3、消费端限流

由于服务器的性能问题,一次处理的请求数有限,如果莫个时刻的请求数极大超过了阈值,服务器可能会崩溃,所以就需要限流,让服务器每次从服务器中拉取一定的请求数,就不会出现宕机 。 【consumers-per-queue:规定一次拉取的请求数】

限流必须要求消费端的签收方式是手动【manual】或根据异常签收【auto】

  • 配置设置

  • 监听类
package com.qiumin.listener;

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


public class AckListener implements ChannelAwareMessageListener {
    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        //接收到消息
        System.out.println(new String(message.getBody()));
       try {
           //处理业务
          // int i=1/0;  制造异常
           int i = 1+1;
           System.out.println("处理业务..."+":"+i);
           //处理正常,手动签收
           Thread.sleep(2000); //方便看到效果,线程睡眠2秒,消费时2秒签收一个
           System.out.println("成功!!!");
           channel.basicAck(deliveryTag,true); //true代表签收多个消息
       }catch (Exception e){
           //处理出现异常拒收
           System.out.println("异常");
           channel.basicNack(deliveryTag,true,true);
           //第一个true代表拒收多个消息,第二个true代表是否将消息重回队列,重新发送
       }

    }
}
4、TTL 过期时间

TTL全称Time to live 存活时间、过期时间

  • 对队列设置过期时间,一段时间后整个队列过期,队列中的消息也过期成为死信。
  • 对消息设置过期时间,一段时间后该消息过期,其他消息无影响,消息过期后只有在队头是才判断过期移除成为死信。
  • 创建队列时设置过期时间。

带有TTL的队列

创建 TTL 队列


    
        
      
            
        
    
    
        
            
        
    
  • 队列会在3秒后过期,队列中的消息也过期,且队列过期后不能在向队列中发送消息

带有TTL的消息 MessagePostProcessor类

 @Test
    public void ttlMessage(){
        //消息后处理对象,设置一些消息的参数 MessagePostProcessor
        MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
            @Override
            public Message postProcessMessage(Message message) throws AmqpException {
                message.getMessageProperties().setExpiration("5000"); //消息5秒后过期
                return message;
            }
        };
        rabbitTemplate.convertAndSend("fanout_exchange","","ttl 带有过期时间的消息",messagePostProcessor);
    }

总结:

  • 设置过期队列使用参数:x-message-ttl
  • 设置过期消息使用参数:expiration
  • 两者都设置了以时间短的为准
5、死信队列

死信队列可以使成为 Dead message后,可以重新发送到另一台交换机,这个交换机称为死信交换机【DLX】

图示:

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

  1. 队列的消息长度达到了限制后再往其中发送消息就会成为死信
  2. 消费者拒绝消费消息,且不把消息放回原队列即 requeue=false
  3. 队列过期,消息还未被消费

配置

 
    
    
        
            
            
            
            
            
            
        
    
    
    
        
            
            
        
    

    
    
    
    
        
            
            
        
    

向普通交换机发送16条消息,由于其绑定的队列最大存储8条消息,所以其他8条消息通过死信交换机转发到死信队列中

  @Test
    public void Test_sx(){
        for (int i = 0; i < 16; i++) {
            rabbitTemplate.convertAndSend("test_exchange_sx","haha.tsx","测试死信队列");
        }
    }

  • 死信交换机和普通的交换机没有区别
  • 消息成为了死信后如果绑定了死信交换机,消息就会被死信交换机路由到死信队列
6、延迟队列

延迟队列,消息延迟消费,相当于定时器,但使用延迟队列效果更好体验更加

  • rabbitmq没有提供延迟队列效果,但可以通过 TTL+死信队列 实现延迟队列。
  • 将需要延迟消费的消息放在ttl队列中,也可是指ttl消息,但该消息在队头时才判断本身就存在延迟不好控制,过期后由死信交换机转发到死信队列,监听死信队列,达到延迟消费的效果。

代码实现

  
    
    
        
            
            
            
            
            
            
        
    
    
    
        
            
            
        
    

    
    
    
    
        
            
            
        
    

监听类

package com.qiumin.listener;

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


public class YCQueueLisener implements ChannelAwareMessageListener {
    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        System.out.println(new String(message.getBody()));
        channel.basicAck(deliveryTag,true);
    }
}

监听死信队列进行消费

    
    

    
    
         
    

  • 普通队列过期后由死信交换机转发到死信队列【延迟队列】,消费端监听延迟队列消费消息
7、日志与监控

rabbitmqctl管理和监控

  • rabbitmqctl list_queues: 查看队列。
  • rabbitmqctl list_exchanges: 查看交换机。
  • rabbitmqctl list_users: 查看用户。
  • rabbitmqctl list_connections: 查看连接。
  • rabbitmqctl list_consumers: 查看消费者信息。
  • rabbitmqctl environment: 查看环境变量。
  • rabbitmqctl list_queues messages_unacknowledged: 查看未被确认的队列。
  • rabbitmqctl list_queues name memory: 查看单个队列内存使用。
  • rabbitmqctl list_queues name messages_ready: 查看准备就绪的队列。

消息追踪 firehose

firrehost的机制是将生产者投递给rabbitmq的消息,rabbitmq投递给消费者的消息按照指定的格式发送到exchange上,这个默认的exchange的名称为amq.rabbitmq.trace,它是一个topic类型的交换机。发送到该交换机的信息的routing key为 publish.exchangename和deliver.queuename,其中第一个为实际的exchange的名称,第二个是实际的queue的名称,分别对应生产者投递到exchange的消息,和消费者从queue上获取的消息。

注意:

  • rabbitmqctl trace_on:开启firehost命令
  • rabbitmqctl trace_off:关闭firehost命令
  • 开启该模式会影响消息的写入功能

可以启动插件:rabbitmq-plugins enable rabbitmq_tracing 比上面的firehost多了一层GUI包装,使用更方便。

RabbitMQ 应用问题 1、消息可靠性保障

消息补偿

  • 两次发送消息,之间有时间间隙,消息消费完消费端会发送确认消息到回调检查服务,服务检测到有确认就把消息的信息存储到数据库,定时检查服务检测该数据库,与消息数据库比对,如果消息数据库多了消息说明有些消息可能丢失,这时就要求生产者重发。

图示:

2、消息幂等性保障

消息的重复消费问题

  • 幂等性指一次和多次某一个资源,对于资源本身应该具有同样的结果,即消息消费消费一次后,再次消费后的结果与消费一次的结果相同。

解决方案 --------- 【乐观锁机制】 :加了乐观锁说明有一个版本号,消费一次完版本号就改变了,其他再次消费时由于版本号不对了就不会被再消费了。

图示:、

qiumin

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

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

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