实验4:Spring Boot 整合 RabbitMQ
实验4:Spring Boot 整合 RabbitMQ
|
一、实验目的和要求
|
|
三、实验完成内容概要(介绍实际完成的内容) Spring Boot 通过 spring-boot-starter-amqp 自动整合 RabbitMQ。核心组件:
|
|
四、实验过程记录(方法、操作截图与步骤说明)
RabbitMQConfig
package com.example.project4.config; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitMQConfig { @Bean public Queue simpleQueue() { // 声明简单模式队列,持久化 return new Queue("simple.queue", true); } @Bean public Queue workQueue() { // 声明工作队列,持久化 return new Queue("work.queue", true); } }
RabbitMQService
package com.example.project4.service; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @Service public class RabbitMQService { @Autowired private RabbitTemplate rabbitTemplate; public void sendSimpleMessage(String message) { rabbitTemplate.convertAndSend("simple.queue", message); System.out.println("消息已发送: " + message); } } 2.创建控制器 RabbitMQController
package com.example.project4.controller; import com.example.project4.service.RabbitMQService; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; @RestController public class RabbitMQController { @Autowired private RabbitMQService rabbitMQService; @GetMapping("/send/simple") public String sendSimpleMessage() { rabbitMQService.sendSimpleMessage("Hello Spring Boot RabbitMQ Simple Mode!"); return "消息已发送"; } } 3.创建消费者 RabbitMQConsumer
package com.example.project4.service; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Service; @Service public class RabbitMQConsumer { private static final Logger log = LoggerFactory.getLogger(RabbitMQConsumer.class); @RabbitListener(queues = "simple.queue") public void receiveSimpleMessage(String message) { log.info("收到简单模式消息: {}", message); } } 3.工作队列模式实践
WorkQueueService
package com.example.project4.service; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @Service public class WorkQueueService { @Autowired private RabbitTemplate rabbitTemplate; public void sendWorkMessage(String message) { rabbitTemplate.convertAndSend("work.queue", message); System.out.println("工作队列消息已发送: " + message); } } 2.创建控制器 WorkQueueController
package com.example.project4.controller; import com.example.project4.service.WorkQueueService; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; @RestController public class WorkQueueController { @Autowired private WorkQueueService workQueueService; @GetMapping("/send/work") public String sendWorkMessage() { for (int i = 0; i < 10; i++) { workQueueService.sendWorkMessage("Work Queue Message " + i); } return "工作队列消息已发送"; } } 3.创建消费者 WorkQueueConsumer
package com.example.project4.service; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Service; @Service public class WorkQueueConsumer { private static final Logger log = LoggerFactory.getLogger(WorkQueueConsumer.class); @RabbitListener(queues = "work.queue") public void receiveWorkMessage(String message) { log.info("收到工作队列消息: {}", message); try { Thread.sleep(1000); // 模拟处理时间 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } Pom.xml <?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.7.15</version> <relativePath/> <!-- lookup parent from repository --> </parent> <groupId>com.example</groupId> <artifactId>project4</artifactId> <version>0.0.1-SNAPSHOT</version> <name>project4</name> <description>project4</description> <properties> <java.version>11</java.version> </properties> <dependencies> <!-- Spring Web --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <!-- RabbitMQ 核心依赖(必须要有) --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <!-- 日志依赖(防止 SLF4J 报错) --> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-api</artifactId> <version>1.7.36</version> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> </project>
|
|
五、结论(结果与分析)
|
实验5:微服务基础组件实践
|
一、实验目的和要求
|
|
三、实验完成内容概要(介绍实际完成的内容)
|
|
四、实验过程记录(方法、操作截图与步骤说明)
spring: rabbitmq: host: localhost port: 5672 username: guest password: guest virtual-host: / # ====================== EUREKA SERVER ====================== --- spring: config: activate: on-profile: eureka application: name: eureka-server server: port: 8761 eureka: client: register-with-eureka: false fetch-registry: false service-url: defaultZone: http://localhost:8761/eureka/ instance: hostname: localhost server: enable-self-preservation: false # ====================== USER SERVICE ====================== --- spring: config: activate: on-profile: user application: name: userservice server: port: 8081 eureka: client: service-url: defaultZone: http://localhost:8761/eureka/ instance: hostname: localhost prefer-ip-address: true logging: level: com.netflix.discovery: OFF org.springframework.cloud: INFO # ====================== ORDER SERVICE ====================== --- spring: config: activate: on-profile: order application: name: orderservice server: port: 8082 eureka: client: service-url: defaultZone: http://localhost:8761/eureka/ instance: prefer-ip-address: true hostname: localhost 1.eurekaserver
package com.example.project5.microservice.eurekaserver; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.cloud.netflix.eureka.server.EnableEurekaServer; @SpringBootApplication @EnableEurekaServer public class eurekaserver { public static void main(String[] args) { System.setProperty("spring.profiles.active", "eureka"); SpringApplication.run(eurekaserver.class, args); } } 步骤 1:启动 eurekaserver 运行 eurekaserver 启动类,激活 eurekaprofile。 访问 http://localhost:8761,确认注册中心启动成功
2.userservice
package com.example.project5.microservice.userservice; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; @RestController public class usercontroller { private final RabbitTemplate rabbitTemplate; public usercontroller(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } @GetMapping("/send") public String sendMessage( @RequestParam String orderId, @RequestParam String userId ) { String msg = "订单ID:" + orderId + ",用户ID:" + userId; rabbitTemplate.convertAndSend( rabbitconfig.ORDER_EXCHANGE, rabbitconfig.ROUTING_KEY, msg ); return "✅ 消息发送成功:" + msg; } }
package com.example.project5.microservice.userservice; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; @RestController public class messageproducer { @Autowired private RabbitTemplate rabbitTemplate; // 这里改成 /sendOrder,避免冲突 @GetMapping("/sendOrder") public String sendMessage( @RequestParam String orderId, @RequestParam String userId) { String message = "订单ID:" + orderId + ",用户ID:" + userId; rabbitTemplate.convertAndSend( rabbitconfig.ORDER_EXCHANGE, rabbitconfig.ROUTING_KEY, message ); return "消息发送成功:" + message; } }
package com.example.project5.microservice.userservice; import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class rabbitconfig { // 全部小写,和你的类名完全一致 public static final String ORDER_QUEUE = "order.queue"; public static final String ORDER_EXCHANGE = "order.exchange"; public static final String ROUTING_KEY = "order.routing.key"; @Bean public Queue orderQueue() { return new Queue(ORDER_QUEUE, true); } @Bean public DirectExchange orderExchange() { return new DirectExchange(ORDER_EXCHANGE); } @Bean public Binding binding(Queue orderQueue, DirectExchange orderExchange) { return BindingBuilder.bind(orderQueue) .to(orderExchange) .with(ROUTING_KEY); } }
package com.example.project5.microservice.userservice; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; @RestController public class usercontroller { private final RabbitTemplate rabbitTemplate; public usercontroller(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } @GetMapping("/send") public String sendMessage( @RequestParam String orderId, @RequestParam String userId ) { String msg = "订单ID:" + orderId + ",用户ID:" + userId; rabbitTemplate.convertAndSend( rabbitconfig.ORDER_EXCHANGE, rabbitconfig.ROUTING_KEY, msg ); return "✅ 消息发送成功:" + msg; } } 运行 userservice 启动类,激活 user profile。 查看 Eureka 控制台,确认 USERSERVICE 状态为 UP。
import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.cloud.netflix.eureka.EnableEurekaClient; @SpringBootApplication @EnableEurekaClient public class orderservice { public static void main(String[] args) { System.setProperty("spring.profiles.active", "order"); SpringApplication.run(orderservice.class, args); } } messageconsumer
import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; @Component public class messageconsumer { @RabbitListener(queues = "order.queue") public void receiveMessage(String message) { System.out.println("========================================"); System.out.println("订单服务收到消息:" + message); System.out.println("========================================"); } } 步骤 3:启动 OrderService(8082) 运行 orderservice 启动类,激活 order profile。 查看 Eureka 控制台,确认 ORDERSERVICE 状态为 UP。
在浏览器访问:http://localhost:8081/send?orderId=1001&userId=2002
Pom.xml <?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.7.18</version> <relativePath/> </parent> <groupId>com.example</groupId> <artifactId>project5</artifactId> <version>0.0.1-SNAPSHOT</version> <properties> <java.version>11</java.version> <spring-cloud.version>2021.0.8</spring-cloud.version> </properties> <dependencies> <!-- Web --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <!-- Eureka Server --> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-netflix-eureka-server</artifactId> </dependency> <!-- Eureka Client + 修复缺失的Jersey依赖 --> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-netflix-eureka-client</artifactId> </dependency> <dependency> <groupId>com.sun.jersey.contribs</groupId> <artifactId>jersey-apache-client4</artifactId> <version>1.19.4</version> </dependency> <!-- RabbitMQ --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> </dependencies> <dependencyManagement> <dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-dependencies</artifactId> <version>${spring-cloud.version}</version> <type>pom</type> <scope>import</scope> </dependency> </dependencies> </dependencyManagement> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> </project> |
实验六:微服务与信息中间件整合实践
|
一、实验目的和要求
|
|
三、实验完成内容概要(介绍实际完成的内容)
|
package com.example.project6.config; import org.springframework.amqp.core.*; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitMQConfig { // 基础订单队列配置 public static final String ORDER_QUEUE = "order.queue"; public static final String ORDER_EXCHANGE = "order.exchange"; public static final String ROUTING_KEY = "order.routing.key"; @Bean public Queue orderQueue() { return QueueBuilder.durable(ORDER_QUEUE).build(); } @Bean public DirectExchange orderExchange() { return ExchangeBuilder.directExchange(ORDER_EXCHANGE).durable(true).build(); } @Bean public Binding orderBinding(Queue orderQueue, DirectExchange orderExchange) { return BindingBuilder.bind(orderQueue).to(orderExchange).with(ROUTING_KEY); } }
import org.springframework.amqp.core.*; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class DeadLetterConfig { // 这里必须加 public static final!!! public static final String PAYMENT_QUEUE = "payment.queue"; public static final String DLX_EXCHANGE = "dlx.exchange"; public static final String DLX_ROUTING_KEY = "dlx.routing.key"; @Bean public Queue paymentQueue() { return QueueBuilder.durable(PAYMENT_QUEUE) .deadLetterExchange(DLX_EXCHANGE) .deadLetterRoutingKey(DLX_ROUTING_KEY) .build(); } @Bean public DirectExchange dlxExchange() { return ExchangeBuilder.directExchange(DLX_EXCHANGE).durable(true).build(); } @Bean public Queue deadLetterQueue() { return QueueBuilder.durable("dead.letter.queue").build(); } @Bean public Binding deadLetterBinding() { return BindingBuilder.bind(deadLetterQueue()).to(dlxExchange()).with(DLX_ROUTING_KEY); } } LazyQueueConfig
package com.example.project6.config; import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.QueueBuilder; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class LazyQueueConfig { public static final String LAZY_QUEUE = "lazy.queue"; @Bean public Queue lazyQueue() { // 声明惰性队列(消息存磁盘,不占内存) return QueueBuilder.durable(LAZY_QUEUE).lazy().build(); } }
import org.springframework.amqp.core.*; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class DelayedMessageConfig { public static final String DELAY_QUEUE = "order.delay.queue"; public static final String DELAY_CONSUME_QUEUE = "order.consume.queue"; public static final String DELAY_EXCHANGE = "order.delay.exchange"; public static final String DEAD_LETTER_EXCHANGE = "order.dlx.exchange"; public static final String DELAY_ROUTING_KEY = "order.delay.routing"; public static final String DEAD_LETTER_ROUTING_KEY = "order.dlx.routing"; @Bean public Queue delayQueue() { return QueueBuilder.durable(DELAY_QUEUE) .ttl(5000) .deadLetterExchange(DEAD_LETTER_EXCHANGE) .deadLetterRoutingKey(DEAD_LETTER_ROUTING_KEY) .build(); } @Bean public Queue consumeQueue() { return QueueBuilder.durable(DELAY_CONSUME_QUEUE).build(); } @Bean public DirectExchange delayExchange() { return new DirectExchange(DELAY_EXCHANGE); } @Bean public DirectExchange delayedDlxExchange() { return new DirectExchange(DEAD_LETTER_EXCHANGE); } @Bean public Binding delayBinding() { return BindingBuilder.bind(delayQueue()).to(delayExchange()).with(DELAY_ROUTING_KEY); } @Bean public Binding delayDlxBinding() { return BindingBuilder.bind(consumeQueue()).to(delayedDlxExchange()).with(DEAD_LETTER_ROUTING_KEY); } } consumer
package com.example.project6.consumer; import com.example.project6.config.DeadLetterConfig; import com.rabbitmq.client.Channel; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; @Component public class PaymentConsumer { // 主队列消费者(模拟业务处理,失败后消息进入死信队列) @RabbitListener(queues = DeadLetterConfig.PAYMENT_QUEUE) public void receivePayment(String msg, Channel channel, Message message) throws IOException { System.out.println("收到支付消息:" + msg); // 模拟业务处理失败,拒绝消息(消息会进入死信队列) channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false); } // 死信队列消费者(处理被拒绝/过期的消息) @RabbitListener(queues = "dead.letter.queue") public void receiveDeadLetter(String msg) { System.out.println("收到死信消息:" + msg); } }
package com.example.project6.consumer; import com.example.project6.config.DelayedMessageConfig; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; @Component public class DelayedMessageConsumer { // 监听真正的消费队列(5秒后收到消息) @RabbitListener(queues = DelayedMessageConfig.DELAY_CONSUME_QUEUE) public void receive(String msg) { System.out.println("✅ 延迟消息已消费:" + msg); } }
package com.example.project6.producer; import com.example.project6.config.RabbitMQConfig; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; @RestController public class OrderProducer { private final RabbitTemplate rabbitTemplate; public OrderProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } @GetMapping("/send/order") public String sendOrderMessage( @RequestParam String orderId, @RequestParam String userId ) { String msg = "订单ID:" + orderId + ",用户ID:" + userId; rabbitTemplate.convertAndSend( RabbitMQConfig.ORDER_EXCHANGE, RabbitMQConfig.ROUTING_KEY, msg ); return "✅ 订单消息发送成功:" + msg; } }
package com.example.project6.producer; import com.example.project6.config.DelayedMessageConfig; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; @RestController public class DelayedMessageProducer { private final RabbitTemplate rabbitTemplate; public DelayedMessageProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } // 发送延迟消息(5秒后收到) @GetMapping("/send/delayed") public String sendDelayed() { String msg = "我是延迟消息,5秒后收到!"; // 发送到延迟队列,5秒后自动进入消费队列 rabbitTemplate.convertAndSend( DelayedMessageConfig.DELAY_EXCHANGE, DelayedMessageConfig.DELAY_ROUTING_KEY, msg ); return "✅ 延迟消息发送成功!延迟 5 秒"; } }
package com.example.project6.service; import com.example.project6.config.RabbitMQConfig; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Service; @Service public class OrderService { private final RabbitTemplate rabbitTemplate; public OrderService(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } public void createOrder(String orderId, String userId) { String msg = "订单创建:ID=" + orderId + ",用户=" + userId; rabbitTemplate.convertAndSend( RabbitMQConfig.ORDER_EXCHANGE, RabbitMQConfig.ROUTING_KEY, msg ); } }
package com.example.project6.service; import com.example.project6.config.DeadLetterConfig; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Service; @Service public class PaymentService { private final RabbitTemplate rabbitTemplate; public PaymentService(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } public void sendPaymentMessage(String orderId) { String msg = "支付请求:订单ID=" + orderId; rabbitTemplate.convertAndSend( DeadLetterConfig.DLX_EXCHANGE, "payment.routing.key", msg ); } } Application.yml server: port: 8080 spring: application: name: project6 rabbitmq: host: localhost port: 5672 username: guest password: guest virtual-host: / publisher-confirm-type: correlated publisher-returns: true listener: simple: acknowledge-mode: manual prefetch: 10 Pom.xml <?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.5.13</version> </parent> <groupId>com.example</groupId> <artifactId>microservice-rabbitmq-demo</artifactId> <version>0.0.1-SNAPSHOT</version> <name>microservice-rabbitmq-demo</name> <properties> <java.version>17</java.version> <spring-cloud.version>2025.0.1</spring-cloud.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-stream-binder-rabbit</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> </dependencies> <dependencyManagement> <dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-dependencies</artifactId> <version>${spring-cloud.version}</version> <type>pom</type> <scope>import</scope> </dependency> </dependencies> </dependencyManagement> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> </project>
2.延迟消息发送成功浏览器返回:
|
|
五、结论(结果与分析)
|
更多推荐











































所有评论(0)