RabbitMQ实战-JavaDemo

news/2024/9/24 8:06:17/

目录

前言

消息生产者

消息消费者

消息确认机制

消息持久化

Maven 依赖

总结


前言

在使用 RabbitMQ 进行消息传递时,了解如何在代码中创建和发布消息(生产者)、接收和处理消息(消费者),以及配置消息确认机制和持久化,是确保系统可靠性和效率的关键。以下将详细解读这些概念,并提供相应的 Java 示例代码和 Maven 依赖。

注:如果是跟着笔者一步步认识RabbitMQ以及手动搭建Rabbit,则需要注意,在搭建RabbitMQ时,笔者运行RabbitMQ使用了指定的vhosts=all,但在本示例中笔者并未指定虚拟主机(vhost),则默认使用/ ,如果你没有创建这个则会报错503,可参考https://blog.csdn.net/StaticKing/article/details/141467806?spm=1001.2014.3001.5502 此篇文章解决,如果你可以正常登录控制台则可在登录控台后,通过页面进行手动添加


消息生产者

消息生产者负责创建并发送消息到 RabbitMQ 服务器的交换机或队列中。

如何在代码中创建和发布消息:

  • 步骤概述:
    1. 创建 ConnectionFactory 并配置 RabbitMQ 服务器的连接参数。
    2. 创建 Connection 并从中创建 Channel
    3. 声明要使用的交换机或队列(如果尚未存在)。
    4. 使用 basicPublish 方法发布消息。

Java 示例:生产者发布消息到队列

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;public class Producer {private final static String QUEUE_NAME = "hello";public static void main(String[] argv) throws Exception {// 1. 创建连接工厂ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("123456");// 2. 创建连接和通道try (Connection connection = factory.newConnection();Channel channel = connection.createChannel()) {// 3. 声明队列(确保队列存在)boolean durable = true; // 队列持久化channel.queueDeclare(QUEUE_NAME, durable, false, false, null);// 4. 创建消息String message = "Hello World!";// 5. 发布消息channel.basicPublish("", QUEUE_NAME,MessageProperties.PERSISTENT_TEXT_PLAIN, // 消息持久化message.getBytes("UTF-8"));System.out.println(" [x] Sent '" + message + "'");}}
}

说明:

  • queueDeclare 方法用于声明一个队列。如果队列已存在且参数匹配,该方法不会有任何效果。
  • basicPublish 方法将消息发送到指定的交换机和路由键。在这个例子中,我们使用默认交换机(空字符串),并将消息直接发布到指定队列。
  • MessageProperties.PERSISTENT_TEXT_PLAIN 设置消息为持久化,以确保在 RabbitMQ 重启后消息不会丢失。

消息消费者

消息消费者负责从 RabbitMQ 服务器的队列中接收并处理消息。

如何在代码中创建和消费消息:

  • 步骤概述:
    1. 创建 ConnectionFactory 并配置连接参数。
    2. 创建 Connection 并从中创建 Channel
    3. 声明要消费的队列(确保队列存在)。
    4. 设置消费回调,定义如何处理接收到的消息。
    5. 使用 basicConsume 方法开始消费。

Java 示例:消费者接收并处理消息

import com.rabbitmq.client.*;public class Consumer {private final static String QUEUE_NAME = "hello";public static void main(String[] argv) throws Exception {// 1. 创建连接工厂ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("123456");// 2. 创建连接和通道Connection connection = factory.newConnection();Channel channel = connection.createChannel();// 3. 声明队列(确保队列存在)boolean durable = true;channel.queueDeclare(QUEUE_NAME, durable, false, false, null);System.out.println(" [*] Waiting for messages. To exit press CTRL+C");// 4. 设置消费者回调DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), "UTF-8");System.out.println(" [x] Received '" + message + "'");try {// 模拟消息处理doWork(message);} finally {System.out.println(" [x] Done");// 5. 手动确认消息channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);}};// 6. 取消自动确认(设置为 false)boolean autoAck = false;channel.basicConsume(QUEUE_NAME, autoAck, deliverCallback, consumerTag -> { });}private static void doWork(String task) {// 模拟耗时操作try {Thread.sleep(1000); // 模拟处理时间} catch (InterruptedException _ignored) {Thread.currentThread().interrupt();}}
}

说明:

  • DeliverCallback 接口用于定义接收到消息时的处理逻辑。
  • basicConsume 方法开始消费消息,参数 autoAck 设置为 false,表示手动确认消息。

消息确认机制

消息确认机制用于确保消息被可靠地处理。RabbitMQ 提供了自动确认和手动确认两种方式,以及在处理失败时的消息重发机制。

确认类型:

  • 自动确认(Auto Acknowledge)

    • 消费者在接收消息后,RabbitMQ 自动将消息标记为已被处理,无需等待消费者的确认。
    • 如果消费者在处理消息时崩溃或失败,消息可能会丢失。
    • 适用于对消息可靠性要求较低的场景。
  • 手动确认(Manual Acknowledge)

    • 消费者在成功处理消息后,显式发送确认信号给 RabbitMQ。
    • 如果消费者处理消息失败或在处理过程中崩溃,RabbitMQ 会重新将消息投递给其他消费者。
    • 提高消息的可靠性,适用于需要确保每条消息都被处理的场景。

Java 示例:手动确认机制

在上述消费者示例中,autoAck 被设置为 false,并在 finally 块中调用了 channel.basicAck 方法进行手动确认。下面是对手动确认机制的更详细说明和示例:

import com.rabbitmq.client.*;public class ManualAckConsumer {private final static String QUEUE_NAME = "hello";public static void main(String[] argv) throws Exception {ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("123456");Connection connection = factory.newConnection();Channel channel = connection.createChannel();boolean durable = true;channel.queueDeclare(QUEUE_NAME, durable, false, false, null);System.out.println(" [*] Waiting for messages. To exit press CTRL+C");DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), "UTF-8");System.out.println(" [x] Received '" + message + "'");try {// 处理消息doWork(message);// 处理成功,发送确认channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);System.out.println(" [x] Done and Acknowledged");} catch (Exception e) {// 处理失败,不发送确认,消息会被重新投递System.out.println(" [x] Error processing message, not acknowledged");}};boolean autoAck = false; // 手动确认channel.basicConsume(QUEUE_NAME, autoAck, deliverCallback, consumerTag -> { });}private static void doWork(String task) throws Exception {// 模拟处理逻辑if (task.contains("fail")) {throw new Exception("Simulated processing failure");}Thread.sleep(1000);}
}

说明:

  • 在消费者处理逻辑中,如果消息处理成功,调用 basicAck 方法发送确认。
  • 如果在处理过程中出现异常,消费者不发送确认,RabbitMQ 将重新投递消息。

消息持久化

为了确保消息在 RabbitMQ 崩溃或重启后不会丢失,需要配置消息持久化。消息持久化包括:

  • 持久化队列:确保队列在服务器重启后仍然存在。
  • 持久化消息:确保消息在服务器重启后不会丢失。

如何实现消息持久化:

  • 持久化队列:在声明队列时,将 durable 参数设置为 true
  • 持久化消息:在发布消息时,将消息属性的 deliveryMode 设置为 2(持久化)。

Java 示例:持久化队列和消息

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;public class PersistentProducer {private final static String QUEUE_NAME = "persistent_task_queue";public static void main(String[] argv) throws Exception {ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("123456");try (Connection connection = factory.newConnection();Channel channel = connection.createChannel()) {// 声明持久化队列boolean durable = true;channel.queueDeclare(QUEUE_NAME, durable, false, false, null);String message = "Persistent Message";// 发布持久化消息channel.basicPublish("", QUEUE_NAME,MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes("UTF-8"));System.out.println(" [x] Sent '" + message + "'");}}
}

说明:

  • channel.queueDeclare 中将 durable 设置为 true,声明持久化队列。
  • 使用 MessageProperties.PERSISTENT_TEXT_PLAIN 设置消息为持久化。

消费者端:

消费者端无需特殊设置来处理持久化消息,只需要正确配置队列即可。

import com.rabbitmq.client.*;public class PersistentConsumer {private final static String QUEUE_NAME = "persistent_task_queue";public static void main(String[] argv) throws Exception {ConnectionFactory factory = new ConnectionFactory();factory.setHost("localhost");factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("123456");Connection connection = factory.newConnection();Channel channel = connection.createChannel();boolean durable = true;channel.queueDeclare(QUEUE_NAME, durable, false, false, null);System.out.println(" [*] Waiting for messages. To exit press CTRL+C");DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), "UTF-8");System.out.println(" [x] Received '" + message + "'");try {doWork(message);channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);} catch (Exception e) {System.out.println(" [x] Error processing message");}};boolean autoAck = false;channel.basicConsume(QUEUE_NAME, autoAck, deliverCallback, consumerTag -> { });}private static void doWork(String task) throws Exception {// 模拟处理Thread.sleep(1000);}
}

总结:

通过配置队列和消息的持久化机制,可以确保即使在 RabbitMQ 崩溃或重启后,消息仍然不会丢失。持久化对于需要高可靠性的消息传递系统尤为重要。


Maven 依赖

为了在 Java 项目中使用 RabbitMQ,需要添加 RabbitMQ Java 客户端库的 Maven 依赖。在 pom.xml 文件中添加以下内容:

<dependencies><dependency><groupId>com.rabbitmq</groupId><artifactId>amqp-client</artifactId><version>5.15.0</version> <!-- 使用最新稳定版本 --></dependency>
</dependencies>

说明:

  • com.rabbitmq:amqp-client 是 RabbitMQ 提供的官方 Java 客户端库。
  • 确保使用最新的稳定版本,以获取最新的功能和安全更新。

总结

理解 RabbitMQ 的消息生产与消费机制,以及相关的确认和持久化机制,对于构建可靠的消息传递系统至关重要。通过正确配置生产者和消费者,并应用手动确认和消息持久化机制,可以确保消息在分布式系统中可靠、安全地传递和处理。结合 Docker 部署 RabbitMQ 和 Maven 集成 Java 客户端,开发者可以快速搭建和测试基于 RabbitMQ 的消息传递解决方案。


http://www.ppmy.cn/news/1517695.html

相关文章

Ubuntu20.04安装 docker和docker-compose环境

Docker简介 Docker 是一个开源的应用容器引擎&#xff0c;它使开发者能够打包他们的应用以及依赖包到一个可移植的容器中&#xff0c;然后发布到任何流行的 Linux 机器上&#xff0c;也可以实现虚拟化。容器是完全使用沙箱机制&#xff0c;相互之间不会有任何接口&#xff08;…

在 CentOS 7 上安装 LNMP 环境:MySQL 8.0、PHP 8.3 和 ThinkPHP 8.0

在 CentOS 7 上安装 LNMP 环境&#xff0c;并配置 MySQL 8.0、PHP 8.3 以及 ThinkPHP 8.0&#xff0c;能够为你的 web 应用程序提供一个强大的开发和运行环境。下面是详细的安装步骤&#xff1a; --- ## 在 CentOS 7 上安装 LNMP 环境&#xff1a;MySQL 8.0、PHP 8.3 和 Thin…

如何处理段错误

在调试代码时&#xff0c;我们会遇到一些状况百出的问题&#xff0c;尤其是段错误&#xff0c;让人头大&#xff1a; 造成段错误的原因主要是内存泄漏&#xff0c;操作空指针&#xff1b; 在很长的代码中&#xff0c;去查找问题是很困难的&#xff0c;这里可以在Linux的ubunt…

Scrapy入门学习

文章目录 Scrapy一. Scrapy简介二. Scrapy的安装1. 进入项目所在目录2. 安装软件包Scrapy3. 验证是否安装成功 三. Scrapy的基础使用1. 创建项目2. 在tutorial/spiders目录下创建保存爬虫代码的项目文件3.运行爬虫4.利用css选择器Scrapy Shell提取数据例如: Scrapy 一. Scrapy…

6个一键生成原创文案实用方法,亲测好用!

在当下的这个自媒体时代&#xff0c;文案创作的需求日益增长。无论是用于社交媒体、广告宣传还是各种内容创作&#xff0c;优质的原创文案都能起到关键作用。但有时候&#xff0c;我们在创作文案的过程中可能会陷入灵感枯竭的困境。但别担心&#xff0c;这里有6个一键生成原创文…

【系统分析师】-缓存

目录 1、常见分类 2、集群切片方式 3、Redis 3.1、分布式存储方式 3.2、数据分片方式 3.3、数据类型 3.4、持久化方案 3.5、内存淘汰机制 3.6、Redis常见问题 4、布隆过滤器 1、常见分类 1、MemCache Memcache是一个高性能的分布式的内存对象缓存系统&#xff0c;用…

Golang 中的 String、rune 和 byte

解释 String Go语言中&#xff0c;string就是只读的采用utf8编码的字节切片(slice) 因此用len函数获取到的长度并不是字符个数&#xff0c;而是字节个数。 for循环遍历输出的也是各个字节。 rune rune是int32的别名&#xff0c;代表字符的Unicode编码&#xff0c;采用4个字…

学习bat脚本

内容包含一些简单命令或小游戏&#xff0c;在乐趣中学习知识。 使用方法&#xff1a; 新建文本文档&#xff0c;将任选其一代码保存到文档中并保存为ASCII编码。将文件后缀改为.bat或.cmd双击运行即可。 一. 关机脚本 1. 直接关机 echo off shutdown -s -t 00秒直接关机。 2…