このトピックでは、rocketmq-spring-boot-starter と rocketmq-v5-client-spring-boot-starter を使用して Spring Boot アプリケーションを ApsaraMQ for RocketMQ 5.0 インスタンスにすばやく接続する方法について説明します。
背景情報
Spring Boot スターターパッケージは、プロデューサーとコンシューマーを作成するためのロジックをカプセル化しています。これらは、Remoting プロトコルの場合は rocketmq-client に、gRPC プロトコルの場合は rocketmq-client-java に依存しています。
その他の SDK については、https://www.alibabacloud.com/help/apsaramq-for-rocketmq/cloud-message-queue-rocketmq-5-x-series/developer-reference/overview-8 をご参照ください。
rocketmq-spring-boot-starter
依存関係の追加
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>{REPLACE_WITH_ACTUAL_VERSION}</version>
</dependency>
-
このスターターは
rocketmq-clientに依存しています。 -
バージョン番号については、Maven リポジトリをご参照ください。
プロジェクトファイルの構成
# メッセージトレース機能
rocketmq.access-channel=CLOUD
rocketmq.name-server={YOUR_ENDPOINT}
rocketmq.consumer.access-key={YOUR_ACCESSKEY_ID}
rocketmq.consumer.secret-key={YOUR_ACCESSKEY_SECRET}
rocketmq.producer.access-key={YOUR_ACCESSKEY_ID}
rocketmq.producer.secret-key={YOUR_ACCESSKEY_SECRET}
#rocketmq.producer.namespaceV2={YOUR_INSTANCE_ID}
rocketmq.producer.group=test
上記のプレースホルダーを実際の値に置き換えてください。中括弧 {} は含めないでください。
-
エンドポイントには、
http://などのプロトコルプレフィックスを追加しないでください。 -
パブリックネットワーク経由でサーバーレスインスタンスに接続する場合は、
namespaceV2パラメーターを構成する必要があります。
メッセージの送信
@Autowired
private RocketMQTemplate rocketMQTemplate;Message<String> msg = MessageBuilder.withPayload("Hello,RocketMQ").build();
// 通常メッセージを送信します。トピックが事前に作成されていることを確認してください。
SendResult sendResult = rocketMQTemplate.syncSend("TEST_TOPIC:mytag", msg); // TEST_TOPIC はトピック名、mytag はタグ名です。
// 遅延メッセージを送信します。トピックを作成する際に、メッセージタイプとして「スケジュール/遅延メッセージ」を選択してください。
SendResult delaySendResult = rocketMQTemplate.syncSendDelayTimeMills("delay:mytag", msg, 6000);
メッセージの消費
@Component
@RocketMQMessageListener(topic = "TEST_TOPIC",selectorExpression = "*", consumerGroup = "GID_test",enableMsgTrace = true,messageModel = MessageModel.CLUSTERING, consumeMode = ConsumeMode.CONCURRENTLY,accessChannel = "CLOUD")
public class MyMQListener implements RocketMQListener<MessageExt> {
@Override
public void onMessage(MessageExt message) {
System.out.println("msg id is " + message.getMsgId() + " , msg body is " + new String(message.getBody()));
}
}
パブリックネットワーク経由でサーバーレスインスタンスに接続する場合は、namespaceV2={YOUR_INSTANCE_ID} 属性を RocketMQMessageListener アノテーションに追加する必要があります。
rocketmq-v5-client-spring-boot-starter
依存関係の追加
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-v5-client-spring-boot-starter</artifactId>
<version>{REPLACE_WITH_ACTUAL_VERSION}</version>
</dependency>
-
このスターターは rocketmq-client-java に依存しています。
-
バージョン番号については、Maven リポジトリをご参照ください。
プロジェクトファイルの構成
# メッセージトレース機能
rocketmq.access-channel=CLOUD
rocketmq.push-consumer.endpoints={YOUR_ENDPOINT}
rocketmq.push-consumer.access-key={YOUR_ACCESSKEY_ID}
rocketmq.push-consumer.secret-key={YOUR_ACCESSKEY_SECRET}
rocketmq.producer.access-key={YOUR_ACCESSKEY_ID}
rocketmq.producer.secret-key={YOUR_ACCESSKEY_SECRET}
rocketmq.producer.namespace={YOUR_INSTANCE_ID}
rocketmq.producer.endpoints={YOUR_ENDPOINT}
メッセージの送信
@Autowired
private RocketMQClientTemplate rocketMQClientTemplate;User user = new User();
user.setName(body);
user.setAge(18);
SendReceipt springbootv5 = rocketMQClientTemplate.syncSendNormalMessage("TEST_TOPIC:mytag", user); // TEST_TOPIC はトピック名、mytag はタグ名です。
メッセージの消費
@Component
@RocketMQMessageListener(topic = "data", consumerGroup = "GID_test", namespace = "{YOUR_INSTANCE_ID}", tag = "*")
public class MyConsumer implements RocketMQListener {
@Override
public ConsumeResult consume(MessageView messageView) {
System.out.println(messageView.getMessageId() + ": body is " + StandardCharsets.UTF_8.decode(messageView.getBody()));
return ConsumeResult.SUCCESS;
}
}
-
サーバーレスインスタンスの場合、
namespace属性にインスタンス ID を構成する必要があります。 -
メッセージボディは
ByteBufferであるため、デコードする必要があります。この例では、StandardCharsets.UTF_8.decode(messageView.getBody())を使用しています。 -
メッセージトレース機能を有効にするには、
enableMsgTrace = trueとaccessChannel = "CLOUD"属性をリスナーのアノテーションに追加してください。
必要な情報の取得
エンドポイントの取得
RocketMQ インスタンスの詳細ページに移動します。基本情報 タブの {type} プロトコルエンドポイント セクションで、エンドポイントとネットワーク情報 を取得し、VPC と インターネット の違いを確認します。
アクセスキー ID とアクセスキーシークレットの取得
インスタンス詳細ページで、左側のナビゲーションペインにあるアクセス制御をクリックし、次にスマート認証をクリックします。[ユーザー名]と[パスワード]フィールドの値は、それぞれ AccessKey ID と AccessKey Secret です。