怎么保证RabbitMq消息不丢失

打印 上一主题 下一主题

主题 878|帖子 878|积分 2634

生产者


  • 做消息的持久化

    • 消息持久化不能保证完全不丢失消息,可以存在存储磁盘的时候还没有存储完成,但是服务宕机了也会导致消息丢失,通过发布确定保证消息不丢失

  • 消息确定机制
  1. import com.rabbitmq.client.*;
  2. public class PersistentProducer {
  3.     private final static String QUEUE_NAME = "persistent_queue";
  4.     public static void main(String[] argv) throws Exception {
  5.         ConnectionFactory factory = new ConnectionFactory();
  6.         factory.setHost("localhost");
  7.         try (Connection connection = factory.newConnection();
  8.              Channel channel = connection.createChannel()) {
  9.             // 声明一个持久化队列
  10.             channel.queueDeclare(QUEUE_NAME, true, false, false, null);
  11.             // 启用生产者确认
  12.             channel.confirmSelect();
  13.             String message = "Persistent message with producer confirm!";
  14.             channel.basicPublish("", QUEUE_NAME,
  15.                 MessageProperties.PERSISTENT_TEXT_PLAIN,
  16.                 message.getBytes());
  17.             // 检查消息是否成功发送
  18.             if (channel.waitForConfirms()) {
  19.                 System.out.println("Message sent successfully!");
  20.             } else {
  21.                 System.out.println("Message failed to send!");
  22.             }
  23.         }
  24.     }
  25. }
复制代码

  • channel.queueDeclare(QUEUE_NAME, true, false, false, null):队列是持久化的,确保 RabbitMQ 重启后队列不会丢失。
  • MessageProperties.PERSISTENT_TEXT_PLAIN:确保消息持久化存储。
  • channel.confirmSelect():启用生产者确认,确保消息乐成送达 RabbitMQ。
  • channel.waitForConfirms():等待确认,如果生产者发送消息时发生失败,会捕获错误。
互换机


  • 选择合适的互换机类型:常见的互换机类型包括 direct、fanout、topic 和 headers,选择正确的类型来确保消息路由正确。
  • 使用死信队列(Dead Letter Exchange, DLX):如果消息因某些原因无法被消耗,可以将消息转发到死信队列进行进一步处理。
  1. import com.rabbitmq.client.*;
  2. public class DirectExchangeProducer {
  3.     private final static String EXCHANGE_NAME = "direct_logs";
  4.     private final static String QUEUE_NAME = "persistent_queue";
  5.     public static void main(String[] argv) throws Exception {
  6.         ConnectionFactory factory = new ConnectionFactory();
  7.         factory.setHost("localhost");
  8.         try (Connection connection = factory.newConnection();
  9.              Channel channel = connection.createChannel()) {
  10.             // 声明交换机和队列
  11.             channel.exchangeDeclare(EXCHANGE_NAME, "direct");
  12.             channel.queueDeclare(QUEUE_NAME, true, false, false, null);
  13.             channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "info");
  14.             String message = "Hello, Direct Exchange!";
  15.             channel.basicPublish(EXCHANGE_NAME, "info",
  16.                 MessageProperties.PERSISTENT_TEXT_PLAIN,
  17.                 message.getBytes());
  18.             System.out.println("Sent: " + message);
  19.         }
  20.     }
  21. }
  22. public class DeadLetterConsumer {
  23.     private final static String DLX_QUEUE = "dlx_queue";
  24.     public static void main(String[] argv) throws Exception {
  25.         ConnectionFactory factory = new ConnectionFactory();
  26.         factory.setHost("localhost");
  27.         try (Connection connection = factory.newConnection();
  28.              Channel channel = connection.createChannel()) {
  29.             // 声明死信队列
  30.             channel.queueDeclare(DLX_QUEUE, true, false, false, null);
  31.             // 创建消费者回调
  32.             DeliverCallback deliverCallback = (consumerTag, delivery) -> {
  33.                 String message = new String(delivery.getBody(), "UTF-8");
  34.                 System.out.println("Dead Letter Queue Received: " + message);
  35.             };
  36.             // 设置死信队列消费者
  37.             channel.basicConsume(DLX_QUEUE, true, deliverCallback, consumerTag -> {});
  38.         }
  39.     }
  40. }
复制代码
消耗者


  • 进行手动应答
  • 消息重试
  1. import com.rabbitmq.client.*;
  2. public class AckConsumer {
  3.     private final static String QUEUE_NAME = "persistent_queue";
  4.     public static void main(String[] argv) throws Exception {
  5.         ConnectionFactory factory = new ConnectionFactory();
  6.         factory.setHost("localhost");
  7.         try (Connection connection = factory.newConnection();
  8.              Channel channel = connection.createChannel()) {
  9.             // 声明一个持久化队列
  10.             channel.queueDeclare(QUEUE_NAME, true, false, false, null);
  11.             // 创建一个消费者回调
  12.             DeliverCallback deliverCallback = (consumerTag, delivery) -> {
  13.                 String message = new String(delivery.getBody(), "UTF-8");
  14.                 System.out.println("Received: " + message);
  15.                 try {
  16.                     // 模拟消息处理
  17.                     if (message.contains("error")) {
  18.                         throw new Exception("Error while processing message");
  19.                     }
  20.                     // 手动确认消息
  21.                     channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
  22.                     System.out.println("Message processed and acknowledged");
  23.                 } catch (Exception e) {
  24.                     // 如果消息处理失败,可以将消息重新放回队列
  25.                     System.out.println("Error processing message, requeueing: " + e.getMessage());
  26.                     channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
  27.                 }
  28.             };
  29.             // 设置手动确认
  30.             channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {});
  31.         }
  32.     }
  33. }
复制代码

  • channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false):消息处理乐成后确认消息。
  • channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true):如果消耗失败,重新将消息投递到队列中,供其他消耗者处理。

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

使用道具 举报

0 个回复

倒序浏览

快速回复

您需要登录后才可以回帖 登录 or 立即注册

本版积分规则

伤心客

金牌会员
这个人很懒什么都没写!

标签云

快速回复 返回顶部 返回列表