【无标题】rabbitmq
实验2:RabbitMQ工作模式实践
|
一、实验目的和要求
|
|
三、实验完成内容概要(介绍实际完成的内容)
|
|
四、实验过程记录(方法、操作截图与步骤说明)
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); } }
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); } }
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); } }
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); } }
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); } }
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> |
|
五、结论(结果与分析)
|
实验3:RabbitMQ 核心特性实践
|
一、实验目的和要求
|
|
三、实验完成内容概要(介绍实际完成的内容)
|
|
四、实验过程记录(方法、操作截图与步骤说明) 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(); } }
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(); } }
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 确认。先运行消费者,再运行生产者。
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(); } } 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(); } }
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)); } }); } } 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(); } } 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(); } }
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. 结论:高优先级消息优先消费,实现业务消息优先级调度。 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. 结论:消息超时未消费自动销毁,避免无效消息堆积。 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(); } }
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. 结论:限制消费者预取数量,实现流量控制,防止消费者被瞬间消息压垮。 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 协议,高可用、数据强一致,适合生产核心业务。 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(); } }
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.集群高可用
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> 五、结论(结果与分析)
|
更多推荐
















































所有评论(0)