如何为Quarkus配置RabbitMQ JMS并实现非阻塞请求回复模型
Quarkus 集成 RabbitMQ JMS 实现请求-回复模型
一、依赖配置
首先添加RabbitMQ JMS客户端及Quarkus JMS扩展依赖:
Maven(pom.xml)
<dependency> <groupId>com.rabbitmq</groupId> <artifactId>amqp-client</artifactId> <version>5.20.0</version> <!-- 使用最新稳定版 --> </dependency> <dependency> <groupId>org.messaginghub</groupId> <artifactId>pooled-jms</artifactId> <version>1.2.5</version> </dependency> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-jms</artifactId> </dependency>
Gradle(build.gradle)
implementation 'com.rabbitmq:amqp-client:5.20.0' implementation 'org.messaginghub:pooled-jms:1.2.5' implementation 'io.quarkus:quarkus-jms'
二、Quarkus 配置(application.properties)
配置RabbitMQ连接及连接池参数,替换为你的实际RabbitMQ信息:
# RabbitMQ JMS 连接基础配置 quarkus.jms.url=tcp://localhost:5672 quarkus.jms.username=guest quarkus.jms.password=guest quarkus.jms.connection-factory-name=ConnectionFactory # 连接池配置(应对高并发请求) quarkus.jms.pool.enabled=true quarkus.jms.pool.max-connections=20 quarkus.jms.pool.idle-timeout=30000
三、请求-回复模型代码实现
1. 请求生产者(非阻塞发送请求)
注入JMS上下文,创建临时队列关联请求与回复,异步等待结果:
import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; import jakarta.jms.JMSContext; import jakarta.jms.Queue; import jakarta.jms.TemporaryQueue; import java.util.concurrent.CompletableFuture; @ApplicationScoped public class RequestProducer { @Inject JMSContext context; @Inject @Named("request-queue") Queue requestQueue; public CompletableFuture<String> sendRequest(String payload) { CompletableFuture<String> future = new CompletableFuture<>(); try { // 创建临时队列接收回复 TemporaryQueue replyQueue = context.createTemporaryQueue(); // 设置消息的回复目的地 var message = context.createTextMessage(payload); message.setJMSReplyTo(replyQueue); // 发送请求 context.createProducer().send(requestQueue, message); // 异步监听回复 context.createConsumer(replyQueue).setMessageListener(msg -> { try { String reply = ((jakarta.jms.TextMessage) msg).getText(); future.complete(reply); } catch (Exception e) { future.completeExceptionally(e); } }); } catch (Exception e) { future.completeExceptionally(e); } return future; } }
2. 请求消费者(处理请求并返回回复)
监听请求队列,处理业务逻辑后将结果发送到指定回复队列:
import jakarta.ejb.MessageDriven; import jakarta.jms.JMSContext; import jakarta.jms.Message; import jakarta.jms.MessageListener; import jakarta.inject.Inject; @MessageDriven(mappedName = "request-queue") public class RequestConsumer implements MessageListener { @Inject JMSContext context; @Override public void onMessage(Message message) { try { String request = ((jakarta.jms.TextMessage) message).getText(); // 模拟业务处理逻辑 String reply = "Processed content: " + request; // 发送回复到请求指定的回复队列 context.createProducer().send(message.getJMSReplyTo(), context.createTextMessage(reply)); } catch (Exception e) { throw new RuntimeException("Failed to process request", e); } } }
3. 同步请求接口(非阻塞对外暴露)
通过REST接口调用生产者,返回CompletionStage实现非阻塞处理:
import jakarta.inject.Inject; import jakarta.ws.rs.GET; import jakarta.ws.rs.Path; import jakarta.ws.rs.QueryParam; import java.util.concurrent.CompletionStage; @Path("/request") public class RequestResource { @Inject RequestProducer producer; @GET public CompletionStage<String> sendRequest(@QueryParam("payload") String payload) { return producer.sendRequest(payload); } }
四、关键注意点
- 非阻塞特性:借助
CompletableFuture和Quarkus的异步支持,避免请求线程阻塞,可高效应对大量并发请求。 - 请求回复匹配:使用临时队列
TemporaryQueue为每个请求绑定唯一回复通道,确保回复精准对应请求。 - 连接池优化:开启JMS连接池可复用连接资源,降低高并发场景下的连接创建开销。
内容的提问来源于stack exchange,提问作者agarwal_achhnera
相关产品推荐
相关产品推荐

