0Pricing
RabbitMQ Messaging & Async Systems · 课时

消息确认与持久化

通过实现手动确认并将消息和队列设为持久化,确保消息可靠性,防止消费者或代理发生故障时数据丢失。

消息确认与持久化 是 CoddyKit 上的免费 RabbitMQ Messaging & Async Systems 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 RabbitMQ Messaging & Async Systems 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 RabbitMQ Messaging & Async Systems 课程共包含 4 节课。

本课时的部分内容尚未翻译,以英文显示。

Introduction to Reliability

In distributed systems, ensuring messages are processed reliably is paramount. What happens if a server crashes? Or if a message consumer fails mid-processing?

This lesson explores two key mechanisms in RabbitMQ to prevent data loss and ensure reliability: Message Acknowledgements and Durability for both messages and queues.

What are Message Acknowledgements?

When a consumer receives a message, it needs to tell RabbitMQ that it has successfully processed it. This confirmation is called an Acknowledgement (or 'ack').

  • Without Acks: If a consumer crashes before processing, the message is lost.
  • With Acks: If a consumer crashes, RabbitMQ knows the message wasn't acknowledged and can redeliver it to another consumer.

Automatic vs. Manual Acknowledgements

RabbitMQ supports two modes for acknowledgements:

  • Automatic (Auto-ack): RabbitMQ considers a message acknowledged as soon as it's delivered to the consumer. This is simple but risky, as messages can be lost if the consumer crashes immediately after receiving but before processing.
  • Manual (Explicit Ack): The consumer explicitly sends an acknowledgement back to RabbitMQ *after* it has successfully processed the message. This is the recommended approach for reliable processing.

Implementing Manual Acknowledgements

Let's see how to implement manual acknowledgements in a Java consumer. We use channel.basicAck() after our message processing logic completes.

Try running this example. The consumer will acknowledge the message after a short delay.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;

public class ConsumerAck {
    private final static String QUEUE_NAME = "ack_queue";

    public static void main(String[] argv) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();

        // Declare a non-durable queue
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

        // Set prefetch count to 1 for fair dispatch (covered in Work Queues)
        channel.basicQos(1);

        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), "UTF-8");
            System.out.println(" [x] Received '" + message + "'");
            try {
                Thread.sleep(1000); // Simulate work
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
            System.out.println(" [x] Done and Acknowledged");
        };
        
        // false means manual acknowledgement
        channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {});

        // Producer to send a message (for testing)
        Channel producerChannel = connection.createChannel();
        producerChannel.queueDeclare(QUEUE_NAME, false, false, false, null);
        producerChannel.basicPublish("", QUEUE_NAME, null, "Hello Ack!".getBytes("UTF-8"));
        System.out.println(" [x] Sent 'Hello Ack!'");
    }
}

Handling Failed Message Processing

What if a consumer fails to process a message? Instead of acknowledging it, you can negatively acknowledge it:

  • channel.basicNack(deliveryTag, multiple, requeue): Rejects one or more messages.
  • channel.basicReject(deliveryTag, requeue): Rejects a single message.

The requeue parameter is crucial. If true, the message is sent back to the queue for another consumer. If false, it's discarded or sent to a Dead Letter Exchange (DLX), which we'll cover in a later lesson.

What is Message Durability?

Acknowledgements handle consumer failures, but what about the RabbitMQ broker itself? If the server crashes or restarts, what happens to messages in the queues?

Message Durability ensures that messages persist on disk and survive a broker restart. This means critical messages are never lost, even if the RabbitMQ server goes down unexpectedly.

Making Messages Persistent

To make a message durable, you need to mark it as 'persistent' when publishing. This tells RabbitMQ to write the message to disk.

We use MessageProperties.PERSISTENT_TEXT_PLAIN for this.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;

public class ProducerPersistent {
    private final static String QUEUE_NAME = "persistent_queue";

    public static void main(String[] argv) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {
            
            // Declare a durable queue first!
            channel.queueDeclare(QUEUE_NAME, true, false, false, null);
            
            String message = "Hello Persistent Message!";
            channel.basicPublish(
                "", 
                QUEUE_NAME, 
                MessageProperties.PERSISTENT_TEXT_PLAIN, // Mark message as persistent
                message.getBytes("UTF-8")
            );
            System.out.println(" [x] Sent '" + message + "'");
        }
    }
}

What is Queue Durability?

Just like messages, queues themselves can be durable. If a queue is not durable, it will be lost if the RabbitMQ broker restarts. Any messages inside it (even persistent ones!) will also be lost.

Therefore, for true reliability, both the queue and the messages within it must be durable.

Declaring a Durable Queue

To make a queue durable, you simply set the durable parameter to true when declaring it. This must be done by both the producer and consumer when they declare the queue.

Run this example. It declares a durable queue and sends a persistent message.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;

public class ProducerDurable {
    private final static String DURABLE_QUEUE = "my_durable_queue";

    public static void main(String[] argv) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {
            
            // Declare a durable queue (durable = true)
            channel.queueDeclare(
                DURABLE_QUEUE, 
                true,  // durable
                false, // exclusive
                false, // autoDelete
                null   // arguments
            );
            
            String message = "Hello from a durable queue!";
            channel.basicPublish(
                "", 
                DURABLE_QUEUE, 
                MessageProperties.PERSISTENT_TEXT_PLAIN, // Persistent message
                message.getBytes("UTF-8")
            );
            System.out.println(" [x] Sent '" + message + "' to durable queue.");
        }
    }
}

Quick Check on Reliability

To ensure maximum message reliability (messages are not lost even if a consumer or broker fails), which combination of features is generally required?

Recap & Next Steps

You've learned how to make your RabbitMQ messaging more reliable!

  • Manual Acknowledgements confirm message processing, preventing loss on consumer failure.
  • Message Durability (persistent messages) ensures messages survive broker restarts.
  • Queue Durability ensures the queue definition itself survives broker restarts.

By combining these, you can build robust systems where messages are rarely lost. Next, you'll explore advanced routing patterns using different types of exchanges!

常见问题解答

「消息确认与持久化」课时是免费的吗?

是的 — 「消息确认与持久化」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 RabbitMQ Messaging & Async Systems 课程的其余内容,请升级到 CoddyKit PRO。 RabbitMQ Messaging & Async Systems 课程共包含 4 节课。

「消息确认与持久化」这节课中我会学到什么?

通过实现手动确认并将消息和队列设为持久化,确保消息可靠性,防止消费者或代理发生故障时数据丢失。 你通过在浏览器中直接运行的动手代码来练习 RabbitMQ Messaging & Async Systems,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 RabbitMQ Messaging & Async Systems 需要有经验吗?

无需任何先前经验。CoddyKit 上的 RabbitMQ Messaging & Async Systems 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。

「消息确认与持久化」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 RabbitMQ Messaging & Async Systems 课中编写并运行代码吗?

能。每节 RabbitMQ Messaging & Async Systems 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. Hello World:简单队列
  2. 工作队列:公平分发
  3. 消息确认与持久化
  4. 使用扇出交换器实现发布/订阅
← 返回 RabbitMQ Messaging & Async Systems