IT评测·应用市场-qidao123.com技术社区

标题: rabbitmq怎样保证消息不丢失 [打印本页]

作者: 守听    时间: 2025-3-26 02:03
标题: rabbitmq怎样保证消息不丢失

一、消息大概丢失的场景

消息丢失大概发生在以下几个阶段:

二、RabbitMQ 消息可靠性保障机制

RabbitMQ 提供了一系列机制,保证消息在各个阶段的可靠通报:
2.1 生产者阶段:消息确认机制(Publisher Confirms)



2.2 RabbitMQ 服务器阶段:持久化机制



2.3 消费者阶段:消息确认机制(Acknowledgments)



2.4 死信队列(Dead Letter Queue,DLQ)



三、RabbitMQ 消息不丢失的实践

以下通过 Java 示例,展示怎样在生产者、RabbitMQ 服务器和消费者阶段分别实现消息的可靠性保障。

四、Java 示例:保证消息不丢失

4.1 环境准备

Maven 依赖
  1. <dependencies>
  2.     <dependency>
  3.         <groupId>com.rabbitmq</groupId>
  4.         <artifactId>amqp-client</artifactId>
  5.         <version>5.17.0</version>
  6.     </dependency>
  7. </dependencies>
复制代码
4.2 生产者:消息确认机制和持久化

  1. import com.rabbitmq.client.*;
  2. public class ReliableProducer {
  3.     private static final String EXCHANGE_NAME = "reliable_exchange";
  4.     private static final String QUEUE_NAME = "reliable_queue";
  5.     private static final String ROUTING_KEY = "reliable_routing_key";
  6.     public static void main(String[] args) throws Exception {
  7.         ConnectionFactory factory = new ConnectionFactory();
  8.         factory.setHost("localhost");
  9.         try (Connection connection = factory.newConnection();
  10.              Channel channel = connection.createChannel()) {
  11.             // 声明交换机和队列
  12.             channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT, true); // 持久化交换机
  13.             channel.queueDeclare(QUEUE_NAME, true, false, false, null); // 持久化队列
  14.             channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
  15.             // 开启消息确认
  16.             channel.confirmSelect();
  17.             String message = "Hello, Reliable RabbitMQ!";
  18.             AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
  19.                     .deliveryMode(2) // 消息持久化
  20.                     .build();
  21.             // 发送消息
  22.             channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, props, message.getBytes());
  23.             System.out.println("Message sent: " + message);
  24.             // 等待确认
  25.             if (channel.waitForConfirms()) {
  26.                 System.out.println("Message confirmed by RabbitMQ");
  27.             } else {
  28.                 System.err.println("Message not confirmed, consider retrying.");
  29.             }
  30.         }
  31.     }
  32. }
复制代码

4.3 消费者:手动确认模式

  1. import com.rabbitmq.client.*;
  2. public class ReliableConsumer {
  3.     private static final String QUEUE_NAME = "reliable_queue";
  4.     public static void main(String[] args) throws Exception {
  5.         ConnectionFactory factory = new ConnectionFactory();
  6.         factory.setHost("localhost");
  7.         try (Connection connection = factory.newConnection();
  8.              Channel channel = connection.createChannel()) {
  9.             // 设置手动消息确认
  10.             boolean autoAck = false;
  11.             // 消费消息
  12.             channel.basicConsume(QUEUE_NAME, autoAck, (consumerTag, delivery) -> {
  13.                 String message = new String(delivery.getBody(), "UTF-8");
  14.                 System.out.println("Received message: " + message);
  15.                 try {
  16.                     // 模拟消息处理
  17.                     processMessage(message);
  18.                     // 消息处理成功,返回 ACK
  19.                     channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
  20.                     System.out.println("Message acknowledged");
  21.                 } catch (Exception e) {
  22.                     // 消息处理失败,拒绝并重新投递
  23.                     channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
  24.                     System.err.println("Message requeued due to processing failure");
  25.                 }
  26.             }, consumerTag -> {});
  27.         }
  28.     }
  29.     private static void processMessage(String message) {
  30.         // 模拟消息处理逻辑
  31.         System.out.println("Processing message: " + message);
  32.     }
  33. }
复制代码

4.4 死信队列配置

  1. import com.rabbitmq.client.*;
  2. import java.util.HashMap;
  3. import java.util.Map;
  4. public class DLQConsumer {
  5.     private static final String MAIN_QUEUE = "main_queue";
  6.     private static final String DLX_EXCHANGE = "dlx_exchange";
  7.     private static final String DLQ = "dlq_queue";
  8.     public static void main(String[] args) throws Exception {
  9.         ConnectionFactory factory = new ConnectionFactory();
  10.         factory.setHost("localhost");
  11.         try (Connection connection = factory.newConnection();
  12.              Channel channel = connection.createChannel()) {
  13.             // 配置死信队列
  14.             channel.exchangeDeclare(DLX_EXCHANGE, BuiltinExchangeType.DIRECT, true);
  15.             channel.queueDeclare(DLQ, true, false, false, null);
  16.             channel.queueBind(DLQ, DLX_EXCHANGE, "dlx_routing_key");
  17.             // 配置主队列与死信交换机绑定
  18.             Map<String, Object> argsMap = new HashMap<>();
  19.             argsMap.put("x-dead-letter-exchange", DLX_EXCHANGE);
  20.             argsMap.put("x-dead-letter-routing-key", "dlx_routing_key");
  21.             channel.queueDeclare(MAIN_QUEUE, true, false, false, argsMap);
  22.             // 消费主队列消息
  23.             channel.basicConsume(MAIN_QUEUE, false, (consumerTag, delivery) -> {
  24.                 String message = new String(delivery.getBody(), "UTF-8");
  25.                 System.out.println("Received message from main queue: " + message);
  26.                 // 模拟处理失败
  27.                 channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
  28.                 System.out.println("Message sent to DLQ");
  29.             }, consumerTag -> {});
  30.         }
  31.     }
  32. }
复制代码

五、总结

通过上述机制和代码示例,RabbitMQ 可以有效保证消息不丢失:

免责声明:如果侵犯了您的权益,请联系站长,我们会及时删除侵权内容,谢谢合作!更多信息从访问主页:qidao123.com:ToB企服之家,中国第一个企服评测及商务社交产业平台。




欢迎光临 IT评测·应用市场-qidao123.com技术社区 (https://dis.qidao123.com/) Powered by Discuz! X3.4