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
相关产品推荐
相关产品推荐

