RabbitMQ Messaging & Async Systems · レッスン

Consumer Acknowledgementsと再キューイング

Consumer acknowledgementsを習得し、処理に失敗したメッセージを再キューイングする方法を学びます。エラーを適切に処理できる、障害耐性の高いコンシューマーを設計します。

レッスン 3/411 ステップ

「Consumer Acknowledgementsと再キューイング」はCoddyKit上の無料RabbitMQ Messaging & Async Systemsレッスンです。 これはレッスン3/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはRabbitMQ Messaging & Async Systems学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 RabbitMQ Messaging & Async Systemsコースには全4レッスンが含まれています。

このレッスンの一部はまだ翻訳されておらず、英語で表示されています。

Reliable Message Consumption

In distributed systems, ensuring messages are processed correctly is vital. What happens if a consumer crashes mid-processing? Or if a message causes an error?

This lesson explores consumer acknowledgements and requeuing, essential techniques for building fault-tolerant message consumers.

ACKs: The Handshake

A consumer acknowledgement (ACK) is a signal sent by the consumer back to RabbitMQ. It tells the broker: "I've successfully received and processed this message."

  • Without an ACK, RabbitMQ assumes the message hasn't been processed.
  • This mechanism prevents message loss if a consumer fails before finishing its work.

Automatic vs. Manual ACKs

RabbitMQ offers two ways to acknowledge messages:

  • Automatic (autoAck=true): Messages are acknowledged immediately upon delivery to the consumer. Simple, but risky if processing fails.
  • Manual (autoAck=false): The consumer explicitly sends an ACK after successful processing. This is the default and recommended for reliable systems.

We'll focus on manual acknowledgements for robustness.

Sending a Message (Producer)

Let's set up a simple Java producer to send a message. This message will be consumed later, and we'll apply manual acknowledgements.

Run this code to send a "Hello RabbitMQ!" message.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;

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

    public static void main(String[] argv) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost"); // Connect to local RabbitMQ

        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {

            channel.queueDeclare(QUEUE_NAME, false, false, false, null);
            String message = "Hello RabbitMQ!";
            channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8));
            System.out.println(" [x] Sent '" + message + "'");
        }
    }
}

Manual Acknowledgements in Action

Now, let's create a consumer that uses manual acknowledgements. Notice the autoAck parameter is set to false when consuming.

The channel.basicAck() call confirms successful processing.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import java.nio.charset.StandardCharsets;

public class MyConsumerAck {
    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();

        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
            System.out.println(" [x] Received '" + message + "'");
            try {
                // Simulate processing work
                Thread.sleep(1000);
                System.out.println(" [x] Done processing '" + message + "'");
                channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); // Manual ACK!
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.err.println(" [!] Processing interrupted.");
                // In a real app, you might re-queue or handle error
            }
        };
        // autoAck is set to false for manual acknowledgements
        channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {});
    }
}

When Processing Fails

What if our consumer encounters an error during message processing? If we only use basicAck, the message is lost even if processing wasn't completed.

RabbitMQ provides ways to inform the broker that a message could not be processed successfully, giving us options to handle it.

Requeuing Failed Messages

If a message fails processing due to a transient error (e.g., database connection down), you might want to retry it later. Use channel.basicNack(deliveryTag, multiple, requeue) or channel.basicReject(deliveryTag, requeue).

  • Setting requeue to true sends the message back to the queue for another consumer to pick up.

Run this consumer. It will intentionally fail twice and then successfully process the message.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.TimeoutException;

public class MyConsumerRequeue {
    private final static String QUEUE_NAME = "ack_queue";
    private static int attemptCount = 0;

    public static void main(String[] argv) throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");

        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();

        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
            long deliveryTag = delivery.getEnvelope().getDeliveryTag();
            System.out.println(" [x] Received '" + message + "' (Attempt: " + (++attemptCount) + ")");

            try {
                if (attemptCount <= 2) { // Simulate failure for first 2 attempts
                    throw new RuntimeException("Simulated processing error!");
                }
                // Simulate successful processing
                Thread.sleep(1000);
                System.out.println(" [x] Successfully processed '" + message + "'");
                channel.basicAck(deliveryTag, false); // ACK on success
            } catch (Exception e) {
                System.err.println(" [!] Error processing '" + message + "': " + e.getMessage());
                channel.basicNack(deliveryTag, false, true); // NACK and requeue!
                System.out.println(" [!] Message '" + message + "' requeued.");
            }
        };
        channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {});
    }
}

Discarding Problematic Messages

Sometimes, a message is fundamentally flawed and will *always* cause an error. Requeuing it repeatedly is pointless and can lead to a "poison message" loop.

In such cases, set requeue to false. This discards the message from the queue. Often, these messages are sent to a Dead Letter Exchange (DLX) for later inspection, but we'll cover DLX in a future lesson.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.TimeoutException;

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

    public static void main(String[] argv) throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");

        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();

        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
            long deliveryTag = delivery.getEnvelope().getDeliveryTag();
            System.out.println(" [x] Received '" + message + "'");

            // Simulate an unrecoverable error
            System.err.println(" [!] Fatal error for '" + message + "'. Discarding.");
            channel.basicNack(deliveryTag, false, false); // NACK and DO NOT requeue!
            System.out.println(" [!] Message '" + message + "' rejected (not requeued).");
        };
        channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {});
    }
}

Idempotent Consumers

When you requeue messages, there's a chance a consumer might process the same message multiple times.

It's crucial to design your consumers to be idempotent. This means processing the same message twice (or more) produces the same result as processing it once, without unintended side effects.

Reliability Check

Consider a RabbitMQ consumer configured with manual acknowledgements. If a message is received but the consumer crashes *before* calling channel.basicAck() or channel.basicNack(), what typically happens to that message?

Recap: Building Reliable Consumers

You've learned how to make your RabbitMQ consumers resilient:

  • Manual Acknowledgements: Give consumers control over message fate.
  • Requeuing: Retry messages for transient failures using basicNack(..., true).
  • Discarding: Prevent poison message loops with basicNack(..., false).
  • Idempotency: Design consumers to handle duplicate deliveries gracefully.

These techniques are fundamental for building robust, fault-tolerant messaging systems.

無料で開始

AI チューターと学ぶ RabbitMQ Messaging & Async Systems — 無料

ブラウザでリアルコードを書いて実行し、24/7 の AI チューターから瞬時にサポートを受け、ウェブまたはアプリで続きから学習できます。

コース
11
レッスン
44

よくある質問

「Consumer Acknowledgementsと再キューイング」レッスンは無料ですか?

はい。「Consumer Acknowledgementsと再キューイング」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、RabbitMQ Messaging & Async Systemsコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 RabbitMQ Messaging & Async Systemsコースには全4レッスンが含まれています。

「Consumer Acknowledgementsと再キューイング」で何を学びますか?

Consumer acknowledgementsを習得し、処理に失敗したメッセージを再キューイングする方法を学びます。エラーを適切に処理できる、障害耐性の高いコンシューマーを設計します。 ブラウザで直接実行するハンズオンコードでRabbitMQ Messaging & Async Systemsを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

RabbitMQ Messaging & Async Systemsを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのRabbitMQ Messaging & Async Systemsは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン3/4です。

「Consumer Acknowledgementsと再キューイング」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このRabbitMQ Messaging & Async Systemsレッスンでコードを書いて実行できますか?

はい。すべてのRabbitMQ Messaging & Async Systemsレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. 永続メッセージと永続キュー
  2. 信頼性を高めるPublisher Confirms
  3. Consumer Acknowledgementsと再キューイング
  4. トランザクションとPublisher Confirmの比較
← RabbitMQ Messaging & Async Systemsに戻る