Spring Boot集成异步RabbitMQ RPC实现请求响应模式问题咨询
基于Spring Boot实现RabbitMQ RPC模式的问题
我正在尝试实现RabbitMQ RPC模式(请求/响应),这对我来说是全新的技术,因此落地过程中遇到了较多困难。
我开发的是基于Spring Boot构建的Web应用,整体架构如下:
- 用户填写携带相关信息的表单后提交,会触发处理控制器,例如
@{/processUser} - 表单提交的信息会封装为对象,发送到RabbitMQ队列
- 响应逻辑在另一个Spring项目服务中执行,该服务接收请求后构建响应并返回
- 响应需要在指定时间内由该Spring项目的其他用户操作完成,超时则返回通用响应
我认为响应端代码需要一直在后台独立线程中等待请求,避免占用Spring Boot应用的主线程,以此实现异步效果。
目前我编写的代码已经可以实现我想要的“异步”效果,但我觉得存在我不了解的更优实现方式,同时我也不清楚多用户访问Web应用时,我使用的匿名线程会出现什么问题。方案不需要完美,只要可用性达标即可 :)
下方代码尚未开发完成,没有实现发送对象、动态生成响应等完整逻辑,仅处于测试阶段。
请求端代码
public String call(String message) throws Exception{ final String corrID = UUID.randomUUID().toString(); String replayQueueName = channel.queueDeclare().getQueue(); AMQP.BasicProperties props = new AMQP.BasicProperties.Builder() .correlationId(corrID) .replyTo(replayQueueName) .build(); channel.basicPublish("", requestQueueName,props,message.getBytes()); final BlockingQueue<String> response = new ArrayBlockingQueue<>(1); String ctag = channel.basicConsume(replayQueueName, true, (consumerTag, delivery) -> { if (delivery.getProperties().getCorrelationId().equals(corrID)) { response.offer(new String(delivery.getBody(), "UTF-8")); } }, consumerTag -> {}); String result = response.take(); channel.basicCancel(ctag); return result; }
该方法在处理控制器中的调用逻辑如下:
try(Connection connection = factory.newConnection()){ channel = connection.createChannel(); System.out.println("Sending request..."); String response = call("Test_Message"); System.out.println(response); }catch (Exception e){ e.printStackTrace(); }
响应端代码
@Bean public ConnectionFactory startFactory(){ return new ConnectionFactory(); } @Bean public Connection startCon(ConnectionFactory factory) throws Exception{ return factory.newConnection(); } @Bean public void reciver(){ new Thread(new Runnable() { @Override public void run() { try{ Channel channel = connection.createChannel(); channel.queueDeclare(RPC_QUEUE_NAME,false,false,false,null); channel.queuePurge(RPC_QUEUE_NAME); channel.basicQos(1); System.out.println("Awaiting rpc requests"); Object monitor = new Object(); DeliverCallback deliverCallback = (consumerTag, delivery) ->{ AMQP.BasicProperties replayProps = new AMQP.BasicProperties.Builder() .correlationId(delivery.getProperties().getCorrelationId()) .build(); String response = "RESPONSE_TESTING"; String message = new String(delivery.getBody(),"UTF-8"); System.out.println(message); channel.basicPublish("",delivery.getProperties().getReplyTo(), replayProps, response.getBytes()); channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); synchronized (monitor){ monitor.notify(); } }; channel.basicConsume(RPC_QUEUE_NAME, false, deliverCallback, (consumerTag -> {})); while(true){ synchronized (monitor){ try{ monitor.wait(); }catch (InterruptedException e){ e.printStackTrace(); } } } }catch (Exception e){ e.printStackTrace(); } } }).start(); }
解决方案
现有实现的问题
- 请求端每次调用都新建Connection和Channel,RabbitMQ的Connection是TCP长连接,频繁创建销毁开销极大,高并发场景下会直接耗尽服务器端口资源。
- 请求端用
response.take()是无限阻塞等待,没有超时机制,一旦响应端故障或者处理超时,请求会一直卡住无法返回。 - 响应端手动创建匿名线程消费,没有异常恢复机制,一旦线程因为MQ连接断开、服务异常等问题终止,就再也无法接收新的请求;同时单例的Connection没有自动重连能力,断连后整个RPC链路直接失效。
- 多用户场景下请求端逻辑是同步阻塞的,会占用Tomcat的请求处理线程,并发量稍高就会占满请求线程池,后续请求直接被拒绝。
低成本优化方案(改动小,可用性达标)
- 请求端优化
- 复用Connection和Channel,直接用Spring Boot官方封装的
RabbitTemplate,内置连接池和RPC调用封装,不用手动写原生客户端逻辑。 - 增加超时机制,将
response.take()替换为response.poll(30, TimeUnit.SECONDS),超时后直接返回预设的通用响应即可。
- 复用Connection和Channel,直接用Spring Boot官方封装的
- 响应端优化
- 不要手动创建线程消费,直接用Spring AMQP的
@RabbitListener注解,框架会自动维护消费线程池、连接自动重连、消息确认逻辑,省去手动写同步锁、循环等待的样板代码,也不会出现线程挂掉无人恢复的问题。 - 如果需要等待其他用户操作完成才返回响应,将请求的correlationId、replyTo队列信息存在本地缓存或者Redis中,等用户操作完成后再根据correlationId找到对应的回调队列,发送响应即可。
- 不要手动创建线程消费,直接用Spring AMQP的
内容的提问来源于stack exchange,提问作者uki34
相关产品推荐
相关产品推荐

