ApsaraMQ for RabbitMQ は、Spring フレームワーク用の SDK を提供しています。このトピックでは、Spring SDK を統合してメッセージを送受信する方法について説明します。
前提条件
-
ApsaraMQ for RabbitMQ コンソールで、インスタンス、仮想ホスト、エクスチェンジ、キューなどのリソースを作成済みであること。詳細については、「ステップ 2: リソースを作成する」をご参照ください。
デモプロジェクト
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 以下。
必要に応じて、次のオプション設定を追加できます。
ステップ 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 の一般的なオプション設定です。
重要な注意事項
クライアントがメッセージ消費に、RabbitTemplate の receiveAndConvert メソッドを使用しているかどうかを確認してください。
このメソッドは、キュー名やタイムアウト期間などのパラメータを受け取ります。タイムアウトが指定されている場合、内部処理ロジックには次の問題があります。メッセージを消費する前に、このメソッドはまず BasicCancel コマンドを送信してコンシューマーのサブスクリプションをキャンセルし、次に Ack を送信してメッセージを確認応答します。
ただし、ApsaraMQ for RabbitMQ が BasicCancel リクエストを受信すると、対応するコンシューマーを即座にキャンセルし、すべての未確認応答メッセージを再キューイングします。これは、クライアントがその後に送信する Ack が実質的に無効であることを意味します。サーバーはそれを有効な確認応答として扱わず、代わりにこれらのメッセージを再配信し、クライアント側で重複消費を引き起こします。
さらに、キューのメッセージスループットが高い場合、頻繁なキャンセルと再キューイング操作によりメッセージの蓄積が発生し、消費の安定性にさらに影響を与える可能性があります。