实验2:RabbitMQ工作模式实践

一、实验目的和要求

  1. 掌握RabbitMQ五种核心工作模式的基本原理与适用场景
  2. 理解Exchange、Queue、Binding等关键组件的作用与关系
  3. 熟悉不同工作模式下消息的分发机制与路由规则
  4. 掌握Java客户端实现各种工作模式的代码编写方法
  5. 通过实践验证各种工作模式的特点与差异

三、实验完成内容概要(介绍实际完成的内容)

  1. 完成简单模式,实现一对一消息收发,验证基础队列模型。
  2. 完成工作队列模式,实现多消费者轮询分发与负载均衡。
  3. 完成发布 / 订阅模式,基于 Fanout 交换机实现消息广播。
  4. 完成路由模式,基于 Direct 交换机实现按路由键精准投递。
  5. 完成通配符模式,基于 Topic 交换机实现 */# 模糊匹配路由。
  6. 对比五种模式的运行机制、分发规则与适用场景,完成实验分析与总结

四、实验过程记录(方法、操作截图与步骤说明)
4.1 简单模式(Simple Mode)
ProducerSimple

package com.atguigu.rabbitmq.Simple;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class ProducerSimple {

    private static final String QUEUE_NAME = "simple_queue";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        try (Connection connection = factory.newConnection();

             Channel channel = connection.createChannel()) {

            channel.queueDeclare(QUEUE_NAME, true, false, false, null);

            String message = "Hello Simple Mode!";

            channel.basicPublish("", QUEUE_NAME, null, message.getBytes());

            System.out.println("消息已发送: " + message);

        }

    }

}

ConsumerSimple

package com.atguigu.rabbitmq.Simple;

import com.rabbitmq.client.*;

import java.io.IOException;

import java.io.UnsupportedEncodingException;

public class ConsumerSimple {

    private static final String QUEUE_NAME = "simple_queue";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        channel.queueDeclare(QUEUE_NAME, true, false, false, null);

        System.out.println("等待接收消息...");

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws UnsupportedEncodingException

            {

                String message = new String(body, "UTF-8");

                System.out.println("收到消息: " + message);

            }

        };

        channel.basicConsume(QUEUE_NAME, true, consumer);

    }

}

  1. 操作步骤
    编写生产者,连接 RabbitMQ 并声明队列,发送消息。
    编写消费者,监听同一队列,接收并打印消息。
    先运行生产者,再运行消费者。
  2. 运行结果生产者:消息已发送: Hello Simple Mode!
    消费者:收到消息: Hello Simple Mode!
    结论:一对一通信,消息仅被消费一次。

4.2 工作队列模式(Work Queue Mode)

WorkQueueProducer

package com.atguigu.rabbitmq.WorkQueue;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class WorkQueueProducer {

    private static final String QUEUE_NAME = "work_queue";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        try (Connection connection = factory.newConnection();

             Channel channel = connection.createChannel()) {

            channel.queueDeclare(QUEUE_NAME, true, false, false, null);

            for (int i = 1; i <= 10; i++) {

                String message = "Task " + i + " (耗时" + (i % 3 + 1) + "秒)";

                channel.basicPublish("", QUEUE_NAME, null, message.getBytes());

                System.out.println("[生产者] 发送任务: " + message);

                Thread.sleep(500);

            }

        }

    }

}

WorkQueueConsumer1

package com.atguigu.rabbitmq.WorkQueue;

import com.rabbitmq.client.*;

import java.io.IOException; // 这里补上!

import java.nio.charset.StandardCharsets;

public class WorkQueueConsumer1 {

    private static final String QUEUE_NAME = "work_queue";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        channel.queueDeclare(QUEUE_NAME, true, false, false, null);

        System.out.println("等待接收工作队列消息...");

        channel.basicQos(1);

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者] 处理任务: " + message);

                try {

                    Thread.sleep(1000 * Integer.parseInt(message.split(" ")[1].replace("(", "").replace(")", "")));

                } catch (InterruptedException e) {

                    Thread.currentThread().interrupt();

                }

                System.out.println("[消费者] 任务完成: " + message);

                channel.basicAck(envelope.getDeliveryTag(), false);

            }

        };

        channel.basicConsume(QUEUE_NAME, false, consumer);

    }

}

WorkQueueConsumer2

package com.atguigu.rabbitmq.WorkQueue;

import com.rabbitmq.client.*;

import java.io.IOException; // 这里补上!

import java.nio.charset.StandardCharsets;

public class WorkQueueConsumer2 {

    private static final String QUEUE_NAME = "work_queue";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        channel.queueDeclare(QUEUE_NAME, true, false, false, null);

        System.out.println("等待接收工作队列消息...");

        channel.basicQos(1);

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者] 处理任务: " + message);

                try {

                    Thread.sleep(1000 * Integer.parseInt(message.split(" ")[1].replace("(", "").replace(")", "")));

                } catch (InterruptedException e) {

                    Thread.currentThread().interrupt();

                }

                System.out.println("[消费者] 任务完成: " + message);

                channel.basicAck(envelope.getDeliveryTag(), false);

            }

        };

        channel.basicConsume(QUEUE_NAME, false, consumer);

    }

}

  1. 操作步骤
    编写生产者循环发送 10 条任务消息。
    编写消费者,设置手动确认与 basicQos(1) 实现公平分发。
    启动两个消费者,再运行生产者。
  2. 运行结果
    两个消费者交替接收消息,实现轮询分发与负载均衡。
    结论:适合任务异步处理,提高系统并发能力。

4.3 发布 / 订阅模式(Publish/Subscribe)

PublishSubscribeProducer

package com.atguigu.rabbitmq.PublishSubscribe;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class PublishSubscribeProducer {

    private static final String EXCHANGE_NAME = "logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        try (Connection connection = factory.newConnection();

             Channel channel = connection.createChannel()) {

            channel.exchangeDeclare(EXCHANGE_NAME, "fanout");

            for (int i = 1; i <= 5; i++) {

                String message = "Message " + i;

                channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes());

                System.out.println("[生产者] 发送消息: " + message);

                Thread.sleep(1000);

            }

        }

    }

}

PublishSubscribeConsumer1

package com.atguigu.rabbitmq.PublishSubscribe;

import com.rabbitmq.client.*;

import java.io.IOException;

import java.nio.charset.StandardCharsets;

public class PublishSubscribeConsumer1 {

    private static final String EXCHANGE_NAME = "logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        // 声明交换机

        channel.exchangeDeclare(EXCHANGE_NAME, "fanout");

        // 声明临时队列(exclusive=true)

        String queueName = channel.queueDeclare("", true, false, true, null).getQueue();

        // 绑定交换机和队列

        channel.queueBind(queueName, EXCHANGE_NAME, "");

        System.out.println("等待接收发布订阅消息...");

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                // ========== 这里修复了编码异常 ==========

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者] 收到消息: " + message);

            }

        };

        channel.basicConsume(queueName, true, consumer);

    }

}

PublishSubscribeConsumer2

package com.atguigu.rabbitmq.PublishSubscribe;

import com.rabbitmq.client.*;

import java.io.IOException;

import java.nio.charset.StandardCharsets;

public class PublishSubscribeConsumer2 {

    private static final String EXCHANGE_NAME = "logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        // 声明交换机

        channel.exchangeDeclare(EXCHANGE_NAME, "fanout");

        // 声明临时队列(exclusive=true)

        String queueName = channel.queueDeclare("", true, false, true, null).getQueue();

        // 绑定交换机和队列

        channel.queueBind(queueName, EXCHANGE_NAME, "");

        System.out.println("等待接收发布订阅消息...");

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                // ========== 这里修复了编码异常 ==========

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者] 收到消息: " + message);

            }

        };

        channel.basicConsume(queueName, true, consumer);

    }

}

  1. 操作步骤
    声明 Fanout 类型交换机,生产者将消息发送到交换机。
    多个消费者创建临时队列并绑定交换机。
    启动多个消费者,运行生产者。
  2. 运行结果
    所有消费者都收到完全相同的全部消息,实现广播效果。
    结论:适用于通知、日志同步等一对多广播场景。

4.4 路由模式(Routing Mode)

RoutingProducer

package com.atguigu.rabbitmq.Routing;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class RoutingProducer {

    private static final String EXCHANGE_NAME = "direct_logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        try (Connection connection = factory.newConnection();

             Channel channel = connection.createChannel()) {

            channel.exchangeDeclare(EXCHANGE_NAME, "direct");

            // 发送不同级别的日志消息

            String[] routingKeys = {"info", "warning", "error"};

            for (String routingKey : routingKeys) {

                String message = "Log message with routing key: " + routingKey;

                channel.basicPublish(EXCHANGE_NAME, routingKey, null, message.getBytes());

                System.out.println("[生产者] 发送消息: " + message);

            }

        }

    }

}

RoutingConsumer1(接info)

package com.atguigu.rabbitmq.Routing;

import com.rabbitmq.client.*;

import java.io.IOException;

import java.nio.charset.StandardCharsets;

public class RoutingConsumer1 {

    private static final String EXCHANGE_NAME = "direct_logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        channel.exchangeDeclare(EXCHANGE_NAME, "direct");

        String queueName = channel.queueDeclare().getQueue();

        channel.queueBind(queueName, EXCHANGE_NAME, "info");

        System.out.println("路由消费者1 [只接收 info] 等待消息...");

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者1-info] 收到:" + message);

            }

        };

        channel.basicConsume(queueName, true, consumer);

    }

}

RoutingConsumer2(接收 error

package com.atguigu.rabbitmq.Routing;

import com.rabbitmq.client.*;

import java.io.IOException;

import java.nio.charset.StandardCharsets;

public class RoutingConsumer2 {

    private static final String EXCHANGE_NAME = "direct_logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        channel.exchangeDeclare(EXCHANGE_NAME, "direct");

        String queueName = channel.queueDeclare().getQueue();

        channel.queueBind(queueName, EXCHANGE_NAME, "error");

        System.out.println("路由消费者2 [只接收 error] 等待消息...");

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者2-error] 收到:" + message);

            }

        };

        channel.basicConsume(queueName, true, consumer);

    }

}

  1. 操作步骤
    声明 Direct 交换机,生产者按 info/warning/error 路由键发送消息。
    消费者绑定指定路由键,只接收匹配消息。
  2. 运行结果
    消费者只接收绑定过的路由键消息,实现精准路由。
    结论:适用于按类型筛选、分级消息场景。

4.5 通配符模式(Topics Mode)

TopicsProducer

package com.atguigu.rabbitmq.Topics;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class TopicsProducer {

    private static final String EXCHANGE_NAME = "topic_logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        try (Connection connection = factory.newConnection();

             Channel channel = connection.createChannel()) {

            channel.exchangeDeclare(EXCHANGE_NAME, "topic");

            // 发送不同主题的消息

            String[] routingKeys = {"stock.usd.eur", "stock.eur.usd", "stock.usd", "stock.eur"};

            for (String routingKey : routingKeys) {

                String message = "Message with routing key: " + routingKey;

                channel.basicPublish(EXCHANGE_NAME, routingKey, null, message.getBytes());

                System.out.println("[生产者] 发送消息: " + message);

            }

        }

    }

}

TopicsConsumer1(只接收 stock.usd. 开头的消息

package com.atguigu.rabbitmq.Topics;

import com.rabbitmq.client.*;

import java.io.IOException;

import java.nio.charset.StandardCharsets;

public class TopicsConsumer1 {

    private static final String EXCHANGE_NAME = "topic_logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        channel.exchangeDeclare(EXCHANGE_NAME, "topic");

        String queueName = channel.queueDeclare().getQueue();

        // 只绑定:stock.usd.xxx 类型消息

        channel.queueBind(queueName, EXCHANGE_NAME, "stock.usd.*");

        System.out.println("【消费者1】等待 stock.usd.* 消息...");

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者1] 收到:" + message + "  key:" + envelope.getRoutingKey());

            }

        };

        channel.basicConsume(queueName, true, consumer);

    }

}

TopicsConsumer2(只接收 .eur 结尾的消息

package com.atguigu.rabbitmq.Topics;

import com.rabbitmq.client.*;

import java.io.IOException;

import java.nio.charset.StandardCharsets;

public class TopicsConsumer2 {

    private static final String EXCHANGE_NAME = "topic_logs";

    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost("127.0.0.1");

        factory.setPort(5672);

        factory.setVirtualHost("/");

        factory.setUsername("guest");

        factory.setPassword("guest");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();

        channel.exchangeDeclare(EXCHANGE_NAME, "topic");

        String queueName = channel.queueDeclare().getQueue();

        // 只绑定:xxx.eur 类型消息

        channel.queueBind(queueName, EXCHANGE_NAME, "*.eur");

        System.out.println("【消费者2】等待 *.eur 消息...");

        DefaultConsumer consumer = new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) throws IOException {

                String message = new String(body, StandardCharsets.UTF_8);

                System.out.println("[消费者2] 收到:" + message + "  key:" + envelope.getRoutingKey());

            }

        };

        channel.basicConsume(queueName, true, consumer);

    }

}

  1. 操作步骤
    声明 Topic 交换机,生产者发送带主题格式的路由键消息。
    消费者使用 *(匹配一个词)、#(匹配多个词)绑定队列。
  2. 运行结果
    消息按模糊匹配规则路由,灵活满足复杂业务分发需求。
    结论:功能最强大,适用于复杂主题订阅场景。

pom.xml

<?xml version="1.0" encoding="UTF-8"?>

<project xmlns="http://maven.apache.org/POM/4.0.0"

         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"

         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">

    <modelVersion>4.0.0</modelVersion>

    <groupId>org.example</groupId>

    <artifactId>module02-work-way</artifactId>

    <version>1.0-SNAPSHOT</version>

    <dependencies>

        <dependency>

            <groupId>com.rabbitmq</groupId>

            <artifactId>amqp-client</artifactId>

            <version>5.20.0</version>

        </dependency>

    </dependencies>

</project>

五、结论(结果与分析)

  1. 简单模式:结构最简单,适合单点消息通知。
  2. 工作队列模式:实现多消费者负载均衡,提升处理效率。
  3. 发布 / 订阅模式:广播消息,所有订阅者全量接收。
  4. 路由模式:精准路由,按规则定向投递。
  5. 通配符模式:模糊匹配,支持复杂主题订阅。

实验3:RabbitMQ 核心特性实践

一、实验目的和要求

  1. 掌握RabbitMQ 4.x版本的核心特性及其应用场景
  2. 理解RabbitMQ 3.x与4.x版本的主要区别
  3. 熟悉RabbitMQ集群、仲裁队列、惰性队列等高级特性的配置与使用
  4. 掌握事务消息、延迟消息、死信队列等机制的实现方法
  5. 通过实践验证RabbitMQ核心特性的功能与效果

三、实验完成内容概要(介绍实际完成的内容)

  1. 完成简单模式,实现一对一消息收发,验证基础队列模型。
  2. 完成工作队列模式,实现多消费者轮询分发与负载均衡。
  3. 完成发布 / 订阅模式,基于 Fanout 交换机实现消息广播。
  4. 完成路由模式,基于 Direct 交换机实现按路由键精准投递。
  5. 完成通配符模式,基于 Topic 交换机实现 */# 模糊匹配路由。
  6. 对比五种模式的运行机制、分发规则与适用场景,完成实验分析与总结

四、实验过程记录(方法、操作截图与步骤说明)
1.公共工具类

RabbitMQUtil

package com.atguigu.rabbitmq.advanced;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class RabbitMQUtil {

    private static final String HOST = "localhost";

    private static final String USERNAME = "guest";

    private static final String PASSWORD = "guest";

    public static Channel getChannel() throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setHost(HOST);

        factory.setUsername(USERNAME);

        factory.setPassword(PASSWORD);

        factory.setVirtualHost("/");

        Connection connection = factory.newConnection();

        return connection.createChannel();

    }

}
RabbitMQClusterUtil

package com.atguigu.rabbitmq.advanced;

import com.rabbitmq.client.Address;

import com.rabbitmq.client.Channel;

import com.rabbitmq.client.Connection;

import com.rabbitmq.client.ConnectionFactory;

public class RabbitMQClusterUtil {

    private static final String USERNAME = "guest";

    private static final String PASSWORD = "guest";

    private static final String VIRTUAL_HOST = "/";

    public static Channel getClusterChannel() throws Exception {

        ConnectionFactory factory = new ConnectionFactory();

        factory.setUsername(USERNAME);

        factory.setPassword(PASSWORD);

        factory.setVirtualHost(VIRTUAL_HOST);

        factory.setAutomaticRecoveryEnabled(true);

        factory.setNetworkRecoveryInterval(5000);

        Address[] addresses = {new Address("localhost", 5672)};

        Connection connection = factory.newConnection(addresses);

        return connection.createChannel();

    }

}

2.消息确认

ConfirmProducer

package com.atguigu.rabbitmq.advanced.confirm;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.Channel;

public class ConfirmProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        channel.queueDeclare("confirm_queue", true, false, false, null);

        channel.confirmSelect();

        String msg = "我是Confirm可靠消息";

        channel.basicPublish("", "confirm_queue", null, msg.getBytes());

        if (channel.waitForConfirms()) {

            System.out.println("【Confirm】消息发送并确认成功");

        }

        channel.close();

    }

}
ConfirmConsumer

package com.atguigu.rabbitmq.advanced.confirm;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.*;

public class ConfirmConsumer {

    public static void main(String[] args) throws Exception {

        // 给 channel 加上 final,解决内部类访问问题

        final Channel channel = RabbitMQUtil.getChannel();

        channel.queueDeclare("confirm_queue", true, false, false, null);

        System.out.println("【Confirm消费者】等待消息...");

        channel.basicConsume("confirm_queue", false, new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag,

                                       Envelope envelope,

                                       AMQP.BasicProperties properties,

                                       byte[] body) {

                System.out.println("收到:" + new String(body));

                try {

                    // 这里就能正常使用 channel 了

                    channel.basicAck(envelope.getDeliveryTag(), false);

                    System.out.println("【手动ACK】确认完成");

                } catch (Exception ignored) {}

            }

        });

    }

}

1. 操作步骤编写生产者,连接 RabbitMQ,开启 Confirm 确认模式,声明队列并发送消息。编写消费者,监听同一队列,关闭自动签收,业务处理完手动 ACK 确认。先运行消费者,再运行生产者。


2. 运行结果生产者:【Confirm】消息发送并确认成功消费者:收到:Confirm 模式可靠消息【手动 ACK】确认完成


3. 结论:通过生产者确认 + 消费者手动签收,保证消息可靠投递、不丢失、不重复消费。


3.事务消息
TransactionProducer

package com.atguigu.rabbitmq.advanced.transaction;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.Channel;

public class TransactionProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        channel.queueDeclare("tx_queue", true, false, false, null);

        try {

            channel.txSelect();

            channel.basicPublish("", "tx_queue", null, "事务消息".getBytes());

            channel.txCommit();

            System.out.println("【事务】提交成功");

        } catch (Exception e) {

            channel.txRollback();

            System.out.println("【事务】回滚成功");

        }

        channel.close();

    }

}
1. 操作步骤编写生产者,连接 RabbitMQ,开启事务模式,声明队列、发送消息,正常提交事务。先运行生产者。

2. 运行结果生产者:【事务】提交成功

3. 结论:RabbitMQ 事务具备原子性,提交消息正常入队,异常可回滚,保证消息不发送。

4.死信队列

DeadLetterProducer

package com.atguigu.rabbitmq.advanced.deadletter;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.Channel;

import java.util.HashMap;

import java.util.Map;

public class DeadLetterProducer {

    public static void main(String[] args) throws Exception {

        final Channel channel = RabbitMQUtil.getChannel();

        // 1. 声明死信交换机和队列

        channel.exchangeDeclare("dead_ex", "direct");

        channel.queueDeclare("dead_queue", true, false, false, null);

        channel.queueBind("dead_queue", "dead_ex", "dead_key");

        // 2. 声明正常队列,并绑定死信配置

        Map<String, Object> queueArgs = new HashMap<String, Object>(); // 去掉菱形<>,兼容Java 5

        queueArgs.put("x-dead-letter-exchange", "dead_ex");

        queueArgs.put("x-dead-letter-routing-key", "dead_key");

        queueArgs.put("x-message-ttl", 5000);

        channel.queueDeclare("normal_queue", true, false, false, queueArgs);

        // 3. 发送消息

        String msg = "我会在5秒后变成死信消息";

        channel.basicPublish("", "normal_queue", null, msg.getBytes());

        System.out.println("【死信生产者】消息已发送,5秒后过期进入死信队列");

        channel.close();

    }

}
DeadLetterConsumer

package com.atguigu.rabbitmq.advanced.deadletter;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.*;

public class DeadLetterConsumer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        // 【修复】消费者也声明一次队列,防止不存在报错

        channel.queueDeclare("dead_queue", true, false, false, null);

        System.out.println("【死信消费者】等待接收死信消息...");

        channel.basicConsume("dead_queue", true, new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag,

                                       Envelope envelope,

                                       AMQP.BasicProperties properties,

                                       byte[] body) {

                System.out.println("【死信队列】收到消息:" + new String(body));

            }

        });

    }

}
1. 操作步骤编写生产者,声明普通队列配置 TTL 和死信交换机绑定。编写消费者监听死信队列。先运行死信消费者,再运行生产者。

2. 运行结果生产者:【死信生产者】消息已发送,5 秒后过期进入死信队列消费者:【死信队列】收到消息:我会在 5 秒后变成死信消息

3. 结论:消息超时、被拒绝会转入死信队列,实现异常消息隔离与事后处理。

5.惰性队列

LazyProducer

package com.atguigu.rabbitmq.advanced.lazy;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.Channel;

import java.util.HashMap;

import java.util.Map;

public class LazyProducer {

    public static void main(String[] args) throws Exception {

        final Channel channel = RabbitMQUtil.getChannel();

        // 兼容 Java 5,不使用菱形语法

        Map<String, Object> queueArgs = new HashMap<String, Object>();

        queueArgs.put("x-queue-mode", "lazy");

        channel.queueDeclare("lazy_queue", true, false, false, queueArgs);

        channel.basicPublish("", "lazy_queue", null, "惰性队列消息(直接落盘)".getBytes());

        System.out.println("【惰性队列】消息已发送");

        channel.close();

    }

}
1. 操作步骤编写生产者,声明队列设置为惰性模式,发送消息。直接运行生产者。

2. 运行结果生产者:【惰性队列】消息已发送

3. 结论:惰性队列消息直接落磁盘,减少内存占用,适合海量消息堆积场景。

6.优先级队列

PriorityProducer

package com.atguigu.rabbitmq.advanced.priority;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.AMQP;

import com.rabbitmq.client.Channel;

import java.util.HashMap;

import java.util.Map;

public class PriorityProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        Map<String, Object> queueArgs = new HashMap<String, Object>();

        queueArgs.put("x-max-priority", 10);

        channel.queueDeclare("priority_queue", true, false, false, queueArgs);

        AMQP.BasicProperties props = new AMQP.BasicProperties.Builder().priority(8).build();

        channel.basicPublish("", "priority_queue", props, "高优先级消息".getBytes());

        System.out.println("【优先级队列】已发送优先级8的消息");

        channel.close();

    }

}
PriorityConsumer

package com.atguigu.rabbitmq.advanced.priority;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.*;

import java.util.HashMap;

import java.util.Map;

public class PriorityConsumer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        // 必须带优先级参数声明

        Map<String, Object> queueArgs = new HashMap<String, Object>();

        queueArgs.put("x-max-priority", 10);

        channel.queueDeclare("priority_queue", true, false, false, queueArgs);

        System.out.println("【优先级消费者】等待消息...");

        channel.basicConsume("priority_queue", true, new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag,

                                       Envelope envelope,

                                       AMQP.BasicProperties properties,

                                       byte[] body) {

                System.out.println("【优先级队列】收到:" + new String(body));

            }

        });

    }

}

1. 操作步骤生产者声明队列设置最大优先级,发送消息携带优先级。消费者监听队列。先运行消费者,再运行生产者。

2. 运行结果生产者:【优先级队列】已发送优先级 8 的消息消费者:【优先级队列】收到:高优先级消息

3. 结论:高优先级消息优先消费,实现业务消息优先级调度。
7.消息超时

TtlProducer

package com.atguigu.rabbitmq.advanced.ttl;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.Channel;

import java.util.HashMap;

import java.util.Map;

public class TtlProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        Map<String, Object> queueArgs = new HashMap<String, Object>();

        queueArgs.put("x-message-ttl", 5000);

        channel.queueDeclare("ttl_queue", true, false, false, queueArgs);

        channel.basicPublish("", "ttl_queue", null, "我5秒后过期".getBytes());

        System.out.println("【TTL超时】消息已发送,5秒后自动删除");

        channel.close();

    }

}

1. 操作步骤生产者声明队列设置全局 TTL 过期时间,发送消息。直接运行生产者。

2. 运行结果生产者:【TTL 超时】消息已发送,5 秒后自动删除

3. 结论:消息超时未消费自动销毁,避免无效消息堆积。
8.预取限流

PrefetchProducer

package com.atguigu.rabbitmq.advanced.prefetch;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.Channel;

public class PrefetchProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        channel.queueDeclare("prefetch_queue", true, false, false, null);

        for (int i = 1; i <= 10; i++) {

            channel.basicPublish("", "prefetch_queue", null, ("消息" + i).getBytes());

        }

        System.out.println("【Prefetch】已发送10条消息");

        channel.close();

    }

}
PrefetchConsumer

package com.atguigu.rabbitmq.advanced.prefetch;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.*;

public class PrefetchConsumer {

    public static void main(String[] args) throws Exception {

        // 关键修复:给 channel 加上 final

        final Channel channel = RabbitMQUtil.getChannel();

        channel.queueDeclare("prefetch_queue", true, false, false, null);

        channel.basicQos(1);

        System.out.println("【Prefetch=1】开启限流消费,一次只处理一条");

        channel.basicConsume("prefetch_queue", false, new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag,

                                       Envelope envelope,

                                       AMQP.BasicProperties properties,

                                       byte[] body) {

                System.out.println("消费:" + new String(body));

                try {

                    Thread.sleep(1000);

                    channel.basicAck(envelope.getDeliveryTag(), false);

                } catch (Exception ignored) {}

            }

        });

    }

}

1. 操作步骤生产者批量发送 10 条消息。消费者设置 basicQos 为 1,手动逐条消费确认。先运行消费者,再运行生产者。

2. 运行结果生产者:【Prefetch】已发送 10 条消息消费者:逐条依次消费每条消息

3. 结论:限制消费者预取数量,实现流量控制,防止消费者被瞬间消息压垮。
9.仲裁队列

QuorumProducer

package com.atguigu.rabbitmq.advanced.quorum;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.Channel;

import java.util.HashMap;

import java.util.Map;

public class QuorumProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        Map<String, Object> queueArgs = new HashMap<String, Object>();

        queueArgs.put("x-queue-type", "quorum");

        channel.queueDeclare("quorum_queue", true, false, false, queueArgs);

        channel.basicPublish("", "quorum_queue", null, "仲裁队列高可靠消息".getBytes());

        System.out.println("【仲裁队列】消息已发送(高可用)");

        channel.close();

    }

}

1. 操作步骤生产者声明仲裁队列并发送消息。直接运行生产者。

2. 运行结果生产者:【仲裁队列】消息已发送(高可用)

3. 结论:基于 Raft 协议,高可用、数据强一致,适合生产核心业务。
10.延迟队列

DelayProducer

package com.atguigu.rabbitmq.advanced.delay;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.AMQP;

import com.rabbitmq.client.Channel;

public class DelayProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        // 声明死信交换机(普通交换机即可实现延迟,不需要插件)

        channel.exchangeDeclare("delay_ex", "direct", true);

        channel.queueDeclare("delay_queue", true, false, false, null);

        channel.queueBind("delay_queue", "delay_ex", "delay_key");

        // 发送消息,设置过期时间10秒

        channel.basicPublish("delay_ex", "delay_key",

                new AMQP.BasicProperties.Builder().expiration("10000").build(),

                "我10秒后到".getBytes());

        System.out.println("【延迟队列】消息已发送(延迟10秒)");

        channel.close();

    }

}
DelayConsumer

package com.atguigu.rabbitmq.advanced.delay;

import com.atguigu.rabbitmq.advanced.RabbitMQUtil;

import com.rabbitmq.client.*;

public class DelayConsumer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQUtil.getChannel();

        channel.queueDeclare("delay_queue", true, false, false, null);

        System.out.println("【延迟消费者】等待延迟消息...");

        channel.basicConsume("delay_queue", true, new DefaultConsumer(channel) {

            @Override

            public void handleDelivery(String consumerTag,

                                       Envelope envelope,

                                       AMQP.BasicProperties properties,

                                       byte[] body) {

                System.out.println("【延迟队列】收到消息:" + new String(body));

            }

        });

    }

}

1. 操作步骤使用延迟交换机,生产者发送带延迟时间的消息。消费者监听延迟队列。先运行消费者,再运行生产者。

2. 运行结果生产者:【延迟队列】消息已发送(延迟 10 秒)消费者:10 秒后收到延迟消息

3. 结论:实现定时延迟投递,可用于订单超时、定时任务场景。

11.集群高可用
ClusterProducer

package com.atguigu.rabbitmq.advanced.cluster;

import com.atguigu.rabbitmq.advanced.RabbitMQClusterUtil;

import com.rabbitmq.client.Channel;

public class ClusterProducer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQClusterUtil.getClusterChannel();

        String queueName = "cluster_queue";

        channel.queueDeclare(queueName, true, false, false, null);

        String msg = "集群高可用可靠消息";

        channel.basicPublish("", queueName, null, msg.getBytes());

        System.out.println("【集群模式】消息已发送到RabbitMQ集群");

        channel.close();

    }

}

ClusterConsumer

package com.atguigu.rabbitmq.advanced.cluster;

import com.atguigu.rabbitmq.advanced.RabbitMQClusterUtil;

import com.rabbitmq.client.*;

public class ClusterConsumer {

    public static void main(String[] args) throws Exception {

        Channel channel = RabbitMQClusterUtil.getClusterChannel();

        channel.queueDeclare("cluster_queue", true, false, false, null);

        System.out.println("【集群消费者】已连接集群,等待消息...");

        channel.basicConsume("cluster_queue", true, new DefaultConsumer(channel){

            @Override

            public void handleDelivery(String consumerTag, Envelope envelope,

                                       AMQP.BasicProperties properties, byte[] body) {

                System.out.println("【集群消费者】已连接集群,成功接收消息:" + new String(body));

            }

        });

    }

}

1. 操作步骤编写集群工具类,开启自动重连与故障恢复。编写集群生产者、消费者。先运行消费者,再运行生产者。

2. 运行结果生产者:【集群模式】消息已发送到 RabbitMQ 集群消费者:【集群消费者】已连接集群,成功接收消息:集群高可用可靠消息

3. 结论:集群支持多节点故障自动转移,自动重连,避免单点故障,保证消息高可用不丢失。

Pom.xml

<?xml version="1.0" encoding="UTF-8"?>

<project xmlns="http://maven.apache.org/POM/4.0.0"

         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"

         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">

    <modelVersion>4.0.0</modelVersion>

    <groupId>org.example</groupId>

    <artifactId>module02-work-way</artifactId>

    <version>1.0-SNAPSHOT</version>

    <depend

encies>

        <dependency>

            <groupId>com.rabbitmq</groupId>

            <artifactId>amqp-client</artifactId>

            <version>5.20.0</version>

        </dependency>

    </dependencies>

</project>

五、结论(结果与分析)

  1. 五种基础消息模型均正常运行,消息分发规则符合预期。
  2. RabbitMQ 4.x 安全性更高,默认禁用匿名用户,界面更现代化。
  3. 仲裁队列基于 Raft 协议,可靠性高于镜像队列,适合高可用场景。
  4. 惰性队列消息直接持久化到磁盘,大幅降低内存占用,适合海量消息堆积。
  5. 事务消息可保证原子性,但性能较低;Confirm 机制更适合高并发。
  6. 死信队列可统一处理异常消息,延迟队列满足定时任务需求。
  7. 手动 ACK + 预取控制可实现公平分发,避免消费者负载不均。

Logo

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

更多推荐