首页技术栈归档照片墙音乐日记随想收藏夹友链留言关于

RabbitMQ入门

写作时间:2026-07-08

RabbitMQ入门

RabbitMQ的角色分类

none:

  • 不能访问management plugin

management:查看自己相关节点信息

  • 列出自己可以通过AMQP登入的虚拟机
  • 查看自己的虚拟机节点 virtual hosts的queues,exchanges和bindings信息
  • 查看和关闭自己的channels和connections
  • 查看有关自己的虚拟机节点virtual hosts的统计信息。包括其他用户在这个节点virtual hosts中的活动信息。

Policymaker

  • 包含management所有权限
  • 查看和创建和删除自己的virtual hosts所属的policies和parameters信息。

Monitoring

  • 包含management所有权限
  • 罗列出所有的virtual hosts,包括不能登录的virtual hosts。
  • 查看其他用户的connections和channels信息
  • 查看节点级别的数据如clustering和memory使用情况
  • 查看所有的virtual hosts的全局统计信息。

Administrator

  • 最高权限
  • 可以创建和删除virtual hosts
  • 可以查看,创建和删除users
  • 查看创建permisssions
  • 关闭所有用户的connections

RabbitMQ入门案例 - Simple 简单模式

1:jdk1.8 2:构建一个maven工程 3:导入rabbitmq的maven依赖 4:启动rabbitmq-server服务 5:定义生产者 6:定义消费者 7:观察消息的在rabbitmq-server服务中的过程

Java原生依赖

<dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.10.0</version>
</dependency>

spring依赖

<dependency>
    <groupId>org.springframework.amqp</groupId>
    <artifactId>spring-amqp</artifactId>
    <version>2.2.5.RELEASE</version>
</dependency>
<dependency>
    <groupId>org.springframework.amqp</groupId>
    <artifactId>spring-rabbit</artifactId>
    <version>2.2.5.RELEASE</version>
</dependency>

springboot依赖

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
systemctl start rabbitmq-server
或者
docker start myrabbit

Producer:服务生成者

package com.yuan.rabbitmq.simple;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;

/**
 * 简单队列生产者
 * @author 111
 */
public class Producer {
    public static void main(String[] args) {
        //1.创建连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接属性
        connectionFactory.setHost("ip地址");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");

        Connection connection=null;
        Channel channel=null;

        try {
            //3.从连接工厂获取
            connection = connectionFactory.newConnection("生产者");
            //4.从连接中获取通道channel
            channel=connection.createChannel();
            //5.声明队列queue存贮消息
            String queueName = "queue1";
            /*
             *  如果队列不存在,则会创建
             *  Rabbitmq不允许创建两个相同的队列名称,否则会报错。
             *  @params1: queue 队列的名称
             *  @params2: durable 队列是否持久化
             *  @params3: exclusive 是否排他,即是否私有的,如果为true,会对当前队列加锁,其他的通道不能访问,并且连接自动关闭
             *  @params4: autoDelete 是否自动删除,当最后一个消费者断开连接之后是否自动删除消息。
             *  @params5: arguments 可以设置队列附加参数,设置队列的有效期,消息的最大长度,队列的消息生命周期等等。
             * */
            channel.queueDeclare(queueName, false, false, false, null);
            //6.准备发送的消息内容
            String message="Hello Word";
            //7.发生消息給中间件rabbitmq-server
            // @params1: 交换机exchange
            // @params2: 队列名称/routing
            // @params3: 属性配置
            // @params4: 发送消息的内容
            channel.basicPublish("",queueName,null,message.getBytes());
            System.out.println("消息发送成功");
        } catch (IOException | TimeoutException e) {
            e.printStackTrace();
        }finally {
            //7.释放连接关闭通道
            if(channel != null&&channel.isOpen()) {
                try {
                    channel.close();
                } catch (IOException | TimeoutException e) {
                    e.printStackTrace();
                }
            }
            if (connection != null ) {
                try {
                    connection.close();
                } catch (IOException e) {
                    e.printStackTrace();
                }
            }
        }
    }
}

Consumer:服务消费者

package com.yuan.rabbitmq.simple;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;

/**
 * 服务消费者
 */
public class Consumer {
    public static void main(String[] args) {
        //1.创建连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接属性
        connectionFactory.setHost("ip地址");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");

        Connection connection=null;
        Channel channel=null;

        try {
            //3.从连接工厂获取
            connection = connectionFactory.newConnection("生产者");
            //4.从连接中获取通道channel
            channel=connection.createChannel();
            //服务消费者,消费  true 自动确认,把MQ中的消息删掉,如果不写就是不把mq中的消息消费掉
            channel.basicConsume("queue1", true, new DeliverCallback() {
                        @Override
                        public void handle(String consumerTag, Delivery message) throws IOException {
                            System.out.println("收到的消息是" + new String(message.getBody(), "utf-8"));
                        }
                        //如果又异常接收失败
                    }, new CancelCallback() {
                        @Override
                        public void handle(String consumerTag) throws IOException {
                            System.out.println("接收失败");
                        }
                    });
            System.out.println("开始接收消息");
            //不关闭阻塞
            System.in.read();

        } catch (IOException | TimeoutException e) {
            e.printStackTrace();
        }finally {
            //7.释放连接关闭通道
            if(channel != null&&channel.isOpen()) {
                try {
                    channel.close();
                } catch (IOException | TimeoutException e) {
                    e.printStackTrace();
                }
            }
            if (connection != null ) {
                try {
                    connection.close();
                } catch (IOException e) {
                    e.printStackTrace();
                }
            }
        }
    }
}

RabbitMQ抽取工具类

package com.yuan.rabbitmq.utils;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
/**
 * 创建工厂的工具类
 */
public class RabbitMqUtils {
    //得到一个连接的channel
    public static Channel getChannel() throws Exception {
        //1.创建连接工厂
        ConnectionFactory connectionFactory = new ConnectionFactory();
        //2.设置连接属性
        connectionFactory.setHost("服务器ip地址");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");
        Connection connection = connectionFactory.newConnection("生产者");
        Channel channel = connection.createChannel();
        return channel;

    }
}

搭建消息生产者

package com.yuan.rabbitmq.two;
import com.rabbitmq.client.Channel;
import com.yuan.rabbitmq.utils.RabbitMqUtils;
import java.util.Scanner;

/**
 * 生产者
 */
public class Task01 {
    //消息队列的名称
    public static  final String QUEUE_NAME="queue1";

    public static void main(String[] args) throws Exception {
        //通过工具类获取连接
      Channel channel = RabbitMqUtils.getChannel();
        /**
         * 生成一个队列
         * 1.队列名称
         * 2.队列里面的消息是否持久化(磁盘)默认情况消息存储在内存中
         * 3.该队列是否只供一个消费者进行消费是否进行消息共享,true可以多个消费者消费false:只能一个消费者消费*
         * 4.是否自动删除最后一个消费者端开连接以后该队一句是否自动删除 true自动删除false不自动删除
         *5.其它参数
         */
        channel.queueDeclare(QUEUE_NAME,false,false,false,null);
        //从控制台接收消息
        Scanner scanner = new Scanner(System.in);
        while (scanner.hasNext()) {
            String message=scanner.next();
            /**
             * 1.发送一个消费
             * 2.路由的key值是哪个  本次是队列的名称
             * 3.其他参数消息
             * 4.发送消息的消息体
             */
            channel.basicPublish("",QUEUE_NAME,null,message.getBytes("utf-8"));
            System.out.println(message+"消息发送成功");
        }
    }
}

搭建消息消费者

package com.yuan.rabbitmq.two;

import com.rabbitmq.client.CancelCallback;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DeliverCallback;
import com.yuan.rabbitmq.utils.RabbitMqUtils;

/**
 * 工作线程,消费者
 */
public class Worker01 {
    //消息队列的名称
    public static  final String QUEUE_NAME="queue1";

    public static void main(String[] args) throws Exception {
        //通过工具类获取连接
        Channel channel = RabbitMqUtils.getChannel();
        //消息的接收
        DeliverCallback deliverCallback=((consumerTag, message) -> {
            System.out.println("接收到的消息"+new String(message.getBody()));
        });
        //消息取消接收收,执行下面的内容
        CancelCallback cancelCallback=(consumerTag -> {
            System.out.println(consumerTag+"消息者取消消费接口回调逻辑");
        });
        /**
         * 水消费者消费消息
         * 1.消费哪个队列
         * 2.消费成功之后是否要自动应答true 代表的自动应答 false 代表手动应答
         * 3.消费者未成功消费的回调
         * 4.消费者取录消费的回调
         */
        System.out.println("等待接收消息.....");
        channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback);

    }
}

消息应答基本介绍

概念

消费者完成一个任务可能需要一段时间,如果其中一个消费者处理一个长的任务并仅只完成了部分突然它挂掉了,会发生什么情况。RabbitMQ一旦向消费者传递了一条消息,便立即将该消息标记为删除。在这种情况下,突然有个消费者挂掉了,我们将丢失正在处理的消息。以及后续发送给该消费这的消息,因为它无法接收到。

为了保证消息在发送过程中不丢失,rabbitmq引入消息应答机制,消息应答就是:消费者在接收到消息并且处理该消息之后,告诉rabbitmq它已经处理了,rabbitmq可以把该消息删除了。

自动应答

消息发送后立即被认为已经传送成功,这种模式需要在高吞吐量和数据传输安全性方面做权衡,因为这种模式如果消息在接收到之前,消费者那边出现连接或者channel 关闭,那么消息就丢失了,当然另一方面这种模式消费者那边可以传递过载的消息,没有对传递的消息数量进行限制,当然这样有可能使得消费者这边由于接收太多还来不及处理的消息,导致这些消息的积压,最终使得内存耗尽,最终这些消费者线程被操作系统杀死,所以这种模式仅适用在消费者可以高效并以某种速率能够处理这些消息的情况下使用。

消息应答的方法

A.Channel.basicAck(用于肯定确认)

RabbitMQ已知道该消息并且成功的处理消息,可以将其丢弃了

B.Channel. basicNack(用于否定确认)

C.Channel. basicReject(用于否定确认)

与Channel.basicNack.相比少一个参数不处理该消息了直接拒绝,可以将其丢弃了

Multiple的解释

手动应答的好处是可以批量应答并且减少网络拥堵

multiple 的 true和 false代表不同意思

true代表批量应答channel上未应答的消息

比如说channel上有传送tag 的消息 5,6,7,8当前tag是8那么此时5-8的这些还未应答的消息都会被确认收到消息应答

false同上面相比

avatar

yuanyourdomain

写代码,做研究,记录生活。

RECOMMENDED

MyBatis 动态 SQL

2026-07-08

Nginx 基础入门

2026-07-08

Maven 多模块与私服

2026-07-08