You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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();
}

解决方案

现有实现的问题

  1. 请求端每次调用都新建Connection和Channel,RabbitMQ的Connection是TCP长连接,频繁创建销毁开销极大,高并发场景下会直接耗尽服务器端口资源。
  2. 请求端用response.take()是无限阻塞等待,没有超时机制,一旦响应端故障或者处理超时,请求会一直卡住无法返回。
  3. 响应端手动创建匿名线程消费,没有异常恢复机制,一旦线程因为MQ连接断开、服务异常等问题终止,就再也无法接收新的请求;同时单例的Connection没有自动重连能力,断连后整个RPC链路直接失效。
  4. 多用户场景下请求端逻辑是同步阻塞的,会占用Tomcat的请求处理线程,并发量稍高就会占满请求线程池,后续请求直接被拒绝。

低成本优化方案(改动小,可用性达标)

  1. 请求端优化
    • 复用Connection和Channel,直接用Spring Boot官方封装的RabbitTemplate,内置连接池和RPC调用封装,不用手动写原生客户端逻辑。
    • 增加超时机制,将response.take()替换为response.poll(30, TimeUnit.SECONDS),超时后直接返回预设的通用响应即可。
  2. 响应端优化
    • 不要手动创建线程消费,直接用Spring AMQP的@RabbitListener注解,框架会自动维护消费线程池、连接自动重连、消息确认逻辑,省去手动写同步锁、循环等待的样板代码,也不会出现线程挂掉无人恢复的问题。
    • 如果需要等待其他用户操作完成才返回响应,将请求的correlationId、replyTo队列信息存在本地缓存或者Redis中,等用户操作完成后再根据correlationId找到对应的回调队列,发送响应即可。

内容的提问来源于stack exchange,提问作者uki34

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 06:39:01