すべてのプロダクト
Search
ドキュメントセンター

ApsaraMQ for RabbitMQ:Spring 統合

最終更新日:Aug 04, 2026

ApsaraMQ for RabbitMQ は、Spring フレームワーク用の SDK を提供しています。このトピックでは、Spring SDK を統合してメッセージを送受信する方法について説明します。

前提条件

デモプロジェクト

SpringBootDemo.zip をクリックして、デモプロジェクトをダウンロードしてください。

ステップ 1: パラメータの設定

application.properties または application.yml ファイルで、設定パラメータを設定します。次の例では、application.properties ファイルを使用します。

# エンドポイント。ApsaraMQ for RabbitMQ コンソールの [Instance Details] ページでエンドポイントを確認します。
spring.rabbitmq.host=XXXXXX.amqp.aliyuncs.com
# ApsaraMQ for RabbitMQ への接続に使用するポート。
spring.rabbitmq.port=5672
# インスタンスの静的ユーザー名。ApsaraMQ for RabbitMQ コンソールの [Static Accounts] ページでユーザー名を確認します。
spring.rabbitmq.username=******
# インスタンスの静的パスワード。ApsaraMQ for RabbitMQ コンソールの [Static Accounts] ページでパスワードを確認します。
spring.rabbitmq.password=******
# 論理的な分離を提供する仮想ホスト。ApsaraMQ for RabbitMQ コンソールの [Vhosts] ページで仮想ホストを確認します。
spring.rabbitmq.virtual-host=test_vhost
# メッセージの確認応答 (Ack) モード。
# 1. none: コンシューマーがメッセージを受信した後、消費の成否にかかわらず、サーバーはメッセージが正常に処理されたと見なします。これは RabbitMQ の autoAck モードです。
# 2. auto: メッセージが正常に消費された後、クライアントは自動的に ack を送信します。メッセージの処理に失敗した場合、クライアントは nack を送信するか、例外をスローします。Channel.basicAck() を明示的に呼び出す必要はありません。
# 3. manual: Ack を手動で送信します。メッセージが正常に消費された後、Channel.basicAck() を明示的に呼び出す必要があります。
spring.rabbitmq.listener.simple.acknowledge-mode=manual

# キャッシュモードを CONNECTION に設定します。ApsaraMQ for RabbitMQ は分散マルチノードアーキテクチャを使用します。CONNECTION モードでは、クライアントはクラスター内の複数のサービスノードにバランスよく接続できます。この方法により、負荷のホットスポットを効果的に防止し、メッセージの送信と消費の効率を向上させることができます。
# 注意: CONNECTION モードでは、RabbitAdmin による自動宣言 (エクスチェンジ、キュー、バインディングの自動作成) は有効になりません。必要なトポロジーリソースを手動で宣言する必要があります。
spring.rabbitmq.cache.connection.mode=connection
# 必要に応じて値を調整してください。
spring.rabbitmq.cache.connection.size=50
# 必要に応じて値を調整してください。
spring.rabbitmq.cache.channel.size=1

# コンシューマーが一度に処理できる未確認応答 (Ack) メッセージの最大数 (QoS)。ApsaraMQ for RabbitMQ サーバーは、min{prefetch, 100} を QoS 値として使用します。コンシューマーの処理能力が低い場合は、この値を減らしてください。
spring.rabbitmq.listener.simple.prefetch=100
# RabbitMQ リスナーの同時実行コンシューマーの最小数。必要に応じて値を調整してください。
spring.rabbitmq.listener.simple.concurrency=10
# RabbitMQ リスナーの同時実行コンシューマーの最大数。消費速度が十分に高い場合、クライアントは max-concurrency のコンシューマーを起動してメッセージを消費します。
spring.rabbitmq.listener.simple.max-concurrency=20

接続とチャネルのベストプラクティス

ApsaraMQ for RabbitMQ は、バックエンドで分散マルチノードアーキテクチャを使用します。サーバー側での単一ノードのパフォーマンスボトルネックを回避するために、クライアント接続を設定する際には次のベストプラクティスに従ってください。

  • CONNECTION モードの使用:クライアントがサーバーへの複数の TCP 接続を確立するように設定します。CONNECTION モードでは、単一のクライアントが複数のバックエンドサービスノードに接続を分散し、クラスター内の複数のノードのリソースを活用します。クライアントあたりの接続数 (CONNECTION) は、クライアントのトラフィック量に基づいて調整し、数十個程度に設定することを推奨します。

  • 頻繁な接続の確立と切断の回避:接続を頻繁に作成および切断しないでください。頻繁な接続の確立と切断の繰り返しは、メッセージの送受信効率を低下させます。操作ごとに接続を作成するのではなく、永続的な接続を維持してください。

  • 接続あたりのチャネル数の制限:単一の接続で作成できるチャネルの数を 50 以下に設定します。単一の接続に多数のチャネルがあると、TCP トラフィックの急激な増加を引き起こす可能性があります。パケットの並び替えが時折発生することと組み合わさると、Linux Netfilter Conntrack メカニズムがパケットを INVALID として分類し、クライアントが RST を送信して接続をリセットする原因となる可能性があります。

次のパラメータを使用して、接続モードと CONNECTION/CHANNEL 数を設定できます。

  • spring.rabbitmq.cache.connection.mode - CONNECTION モードの場合は connection に設定します。

  • spring.rabbitmq.cache.connection.size - 接続数。推奨: 数十個程度。

  • spring.rabbitmq.cache.channel.size - 接続あたりのチャネル数。推奨: 50 以下。

必要に応じて、次のオプション設定を追加できます。

オプション設定

プロパティ

説明

spring.rabbitmq.addresses

クライアントが接続するサーバーアドレス。複数のアドレスはカンマ (,) で区切ります。

spring.rabbitmq.hostspring.rabbitmq.addresses の両方を設定した場合、spring.rabbitmq.addresses パラメータが優先されます。

spring.rabbitmq.dynamic

AmqpAdmin Bean を作成するかどうかを指定します。デフォルト値は true です。

spring.rabbitmq.connection-timeout

接続タイムアウト。単位: ミリ秒。値 0 はタイムアウトなしを示します。

spring.rabbitmq.requested-heartbeat

ハートビートタイムアウト。単位: 秒。デフォルト値は 60 です。

spring.rabbitmq.publisher-confirms

Publisher Confirm メカニズムを有効にするかどうかを指定します。

spring.rabbitmq.publisher-returns

Publisher Return メカニズムを有効にするかどうかを指定します。

spring.rabbitmq.ssl.enabled

SSL 証明書認証を有効にするかどうかを指定します。

spring.rabbitmq.ssl.key-store

SSL 証明書を保持するキーストアへのパス。

spring.rabbitmq.ssl.key-store-password

キーストアにアクセスするためのパスワード。

spring.rabbitmq.ssl.trust-store

信頼できる証明書を含むトラストストアの場所。

spring.rabbitmq.ssl.trust-store-password

トラストストアにアクセスするためのパスワード。

spring.rabbitmq.ssl.algorithm

SSL で使用されるアルゴリズム (例: TLSv1.2)。

spring.rabbitmq.ssl.validate-server-certificate

サーバー証明書の検証を有効にするかどうかを指定します。

spring.rabbitmq.ssl.verify-hostname

ホスト検証を有効にするかどうかを指定します。

spring.rabbitmq.cache.channel.size

キャッシュに保持するチャネルの数。

spring.rabbitmq.cache.channel.checkout-timeout

キャッシュサイズに達したときに、キャッシュからチャネルを取得するためのタイムアウト。

単位: ミリ秒。値 0 は、常に新しいチャネルが作成されることを意味します。

spring.rabbitmq.cache.connection.size

キャッシュされた接続の数。このパラメータは CONNECTION モードでのみ有効です。

spring.rabbitmq.cache.connection.mode

接続キャッシュモード。有効な値は次のとおりです。

  • CHANNEL:CHANNEL モードでは、クライアントは単一の TCP 接続を再利用し、すべてのチャネルがその上で多重化されます。低トラフィックシナリオに適しています。

  • CONNECTION:CONNECTION モードでは、クライアントはクラスター内の複数のサービスノードに分散した複数の TCP 接続を確立します。マルチノードリソースを活用するため、本番環境での使用を推奨します。クライアントあたりの接続数は、数十個程度に設定することを推奨します。

spring.rabbitmq.listener.type

リスナーコンテナのタイプ。有効な値は次のとおりです。

  • simple

  • direct

デフォルト値は simple です。

spring.rabbitmq.listener.simple.auto-startup

アプリケーションの起動時にコンテナを自動的に起動するかどうかを指定します。デフォルト値は true です。

spring.rabbitmq.listener.simple.acknowledge-mode

メッセージの確認応答モード。有効な値は次のとおりです。

  • none:コンシューマーがメッセージを受信した後、消費の成否にかかわらず、サーバーはメッセージが正常に処理されたと見なします。これは RabbitMQ の autoAck モードです。

  • manual:手動で Ack を送信します。メッセージが正常に消費された後、Basic.ack を明示的に呼び出す必要があります。

  • auto:メッセージが正常に消費された後、クライアントは自動的に ack を送信します。メッセージの処理に失敗した場合、クライアントは nack を送信するか、例外をスローします。Basic.ack を明示的に呼び出す必要はありません。

デフォルト値は auto です。

spring.rabbitmq.listener.simple.concurrency

コンシューマーの最小数。

spring.rabbitmq.listener.simple.max-concurrency

コンシューマーの最大数。

spring.rabbitmq.listener.simple.prefetch

コンシューマーが一度に処理できる未確認応答 (Ack) メッセージの最大数。これは、Basic.qos メソッドを呼び出して QoS 値を設定することと同等です。トランザクションを使用する場合、この値はトランザクションサイズ以上である必要があります。

spring.rabbitmq.listener.simple.transaction-size

トランザクションで処理するメッセージの数。

spring.rabbitmq.listener.simple.default-requeue-rejected

拒否されたメッセージを再キューイングするかどうかを指定します。デフォルト値は true です。

spring.rabbitmq.listener.simple.missing-queues-fatal

宣言されたキューがブローカーで使用できない場合にリスナーが失敗するか、または実行時に 1 つ以上のキューが削除された場合にコンテナが停止するかを指定します。デフォルト値は true です。

spring.rabbitmq.listener.simple.idle-event-interval

アイドルコンテナイベントを公開する間隔。単位: ミリ秒。

spring.rabbitmq.template.mandatory

必須メッセージングを有効にするかどうかを指定します。デフォルト値は false です。

spring.rabbitmq.template.receive-timeout

receive() 操作のタイムアウト。

spring.rabbitmq.template.reply-timeout

sendAndReceive() 操作のタイムアウト。

ステップ 2: SDK を使用したメッセージの送受信

メッセージの送信

RabbitMQService で、依存性注入を使用して RabbitTemplate を取得し、その send メソッドを呼び出してメッセージを送信します。

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import java.nio.charset.StandardCharsets;
import java.util.UUID;

@Service
public class RabbitMQService {
    
    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void sendMessage(String exchange, String routingKey, String content) {
        // MessageId を設定します。
        String msgId = UUID.randomUUID().toString();
        MessageProperties messageProperties = new MessageProperties();
        messageProperties.setMessageId(msgId);
        // Message を作成します。
        Message message = new Message(content.getBytes(StandardCharsets.UTF_8), messageProperties);
        /*
         * send() メソッドを呼び出してメッセージを送信します。
         * exchange: エクスチェンジの名前。
         * routingKey: ルーティングキー。
         * message: メッセージの内容。
         * correlationData は Publisher Confirm に使用されます。
         */
        rabbitTemplate.send(exchange, routingKey, message, null);
    }
}

メッセージの消費

@RabbitListener アノテーションを使用してメッセージを消費します。

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

import java.util.Arrays;

@Component
public class MessageListener {

     /**
     * メッセージを受信します。
     * @param message メッセージ。
     * @param channel チャネル。
     * @throws IOException
     * queues を、作成したキュー名に置き換えてください。
     */
    @RabbitListener(queues = "myQueue")
    public void receiveFromMyQueue(Message message, Channel channel) throws IOException {
        // ここにメッセージを消費するビジネスロジックを記述します。
        ...
        // Ack の有効期間 (消費タイムアウト) 内に Ack を返す必要があります。そうしないと、確認は無効になり、メッセージが再配信されます。
        channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
    }
}

以下は、RabbitListener の一般的なオプション設定です。

オプション設定

プロパティ

説明

ackMode

カスタムメッセージ確認応答モード。これは spring.rabbitmq.listener.simple.acknowledge-mode 設定を上書きします。

admin

AMQP リソース管理のための AmqpAdmin への参照。

autoStartup

アプリケーションの起動時にコンテナを自動的に起動するかどうかを指定します。これは spring.rabbitmq.listener.simple.auto-startup 設定を上書きします。

bindings

キューとエクスチェンジ間のバインディングの配列。バインディング情報が含まれます。

concurrency

リスナーコンテナの同時実行スレッド数を設定します。

errorHandler

リスナーメソッドによってスローされた例外のエラーハンドラーを設定します。

exclusive

キューの排他モードを有効にします。これは、コンシューマーがキューへの排他的なアクセス権を持ち、他のコンシューマーがキューからメッセージを受信できないことを意味します。この機能は現在サポートされていません。

queues

このリスナーが待機するキューを宣言します。

queuesToDeclare

明示的に宣言するキューを指定します。

重要な注意事項

クライアントがメッセージ消費に、RabbitTemplatereceiveAndConvert メソッドを使用しているかどうかを確認してください。

このメソッドは、キュー名やタイムアウト期間などのパラメータを受け取ります。タイムアウトが指定されている場合、内部処理ロジックには次の問題があります。メッセージを消費する前に、このメソッドはまず BasicCancel コマンドを送信してコンシューマーのサブスクリプションをキャンセルし、次に Ack を送信してメッセージを確認応答します。

ただし、ApsaraMQ for RabbitMQ が BasicCancel リクエストを受信すると、対応するコンシューマーを即座にキャンセルし、すべての未確認応答メッセージを再キューイングします。これは、クライアントがその後に送信する Ack が実質的に無効であることを意味します。サーバーはそれを有効な確認応答として扱わず、代わりにこれらのメッセージを再配信し、クライアント側で重複消費を引き起こします。

さらに、キューのメッセージスループットが高い場合、頻繁なキャンセルと再キューイング操作によりメッセージの蓄積が発生し、消費の安定性にさらに影響を与える可能性があります。