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

Spring Boot RabbitMQ消费者兼生产者改造问题求助

解决方案

要实现消费者接收消息后向response_queue发送响应的功能,只需对现有消费者代码做以下修改:

1. 注入RabbitTemplate

在SpringConsumer类中注入RabbitTemplate——这是Spring AMQP提供的消息发送工具,用于调用convertAndSend方法发送消息。

2. 修正队列名称拼写错误

你的代码中把目标队列名写成了resonse_queue,正确拼写应为response_queue,需修正常量定义。

3. 添加消息发送逻辑

在cool字段为true的分支中,调用RabbitTemplate的convertAndSend方法发送构造好的响应JSON。

修改后的消费者完整代码

import org.json.JSONException;
import org.json.JSONObject;
import org.springframework.amqp.rabbit.annotation.EnableRabbit;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.web.server.ConfigurableWebServerFactory;
import org.springframework.boot.web.server.WebServerFactoryCustomizer;
import org.springframework.context.annotation.Bean;

@SpringBootApplication
@EnableRabbit
public class SpringConsumer {

    private static final String QUEUE_NAME = "spring-boot4";
    // 修正队列名称拼写
    private static final String RESPONSE_QUEUE_NAME = "response_queue";

    // 注入RabbitTemplate
    @Autowired
    private RabbitTemplate rabbitTemplate;

    public static void main(String[] args) {
        SpringApplication.run(SpringConsumer.class, args);
    }


    @RabbitListener(queues = QUEUE_NAME)
    public void receiveMessage(String message) {
        try {
            JSONObject json = new JSONObject(message);
            if ((boolean) json.get("cool")) {
                System.out.println("Neue 'wahre' Nachricht empfangen: " + json.toString());
                JSONObject responseJson = new JSONObject();
                responseJson.put("response", "success");
                // 发送消息到response_queue
                rabbitTemplate.convertAndSend(RESPONSE_QUEUE_NAME, responseJson.toString());
            } else {
                System.out.println("Neue 'falsche' Nachricht empfangen: " + json.toString());
            }

        } catch (JSONException e) {
            e.printStackTrace();
        }
    }
    

    @Bean
    public WebServerFactoryCustomizer<ConfigurableWebServerFactory> webServerFactoryCustomizer() {
        return factory -> factory.setPort(9090); // 设置端口为9090
    }
}

4. 确保response_queue存在

你需要保证response_queue已经在RabbitMQ中存在,有两种方式:

  • 手动创建:登录RabbitMQ管理控制台,直接创建名为response_queue的队列;
  • 代码声明:如果希望通过代码自动创建队列,可以在消费者或生产者类中添加队列的Bean定义:
@Bean
Queue responseQueue() {
    // durable设为true表示队列持久化
    return new Queue(RESPONSE_QUEUE_NAME, true);
}

如果需要将该队列绑定到交换机(比如生产者中的spring-boot-exchange),还可以添加Binding Bean:

@Bean
Binding responseBinding(Queue responseQueue, TopicExchange exchange) {
    return BindingBuilder.bind(responseQueue).to(exchange).with("response.#");
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 13:05:58