实验4:Spring Boot 整合 RabbitMQ

一、实验目的和要求

  1. 掌握 Spring Boot 与 RabbitMQ 的整合方法与配置流程。
  2. 理解 Spring AMQP 的工作机制与核心组件作用。
  3. 掌握消息生产者、消费者的编写与消息确认机制。
  4. 实现简单队列、工作队列两种经典消息模式。
  5. 学会使用 RabbitTemplate 与 @RabbitListener 完成消息收发。

三、实验完成内容概要(介绍实际完成的内容)

Spring Boot 通过 spring-boot-starter-amqp 自动整合 RabbitMQ。核心组件:

  1. RabbitTemplate:发送消息
  2. @RabbitListener:监听并消费消息
  3. ConnectionFactory:建立连接
  4. Queue/Exchange/Binding:构建消息路由结构

四、实验过程记录(方法、操作截图与步骤说明)

  1. 配置类实现

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

    }

}

  1. 简单模式实践
  1. 创建服务类

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.工作队列模式实践

  1. 创建服务类

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>

五、结论(结果与分析)

  1. 简单模式:实现了点对点的消息通信,生产者发送的消息被消费者成功接收并处理,验证了消息队列的基础收发功能。
  2. 工作队列模式:多个消费者能够监听同一个队列,实现了任务的分发处理,验证了消息队列在异步任务处理中的应用。
  3. 通过本次实验,成功掌握了 Spring Boot 整合 RabbitMQ 的开发流程,理解了消息队列在系统解耦、异步处理中的作用。

实验5:微服务基础组件实践

一、实验目的和要求

  1. 理解微服务架构的核心思想,掌握Spring Cloud与Spring Boot的兼容适配逻辑,明确Spring Cloud 2021.0.x与Spring Boot 2.7.x的技术协同要点。
  2. 掌握微服务核心基础组件(服务注册与发现、配置中心、消息队列)的工作原理,能基于指定技术栈搭建基础微服务架构。
  3. 熟练整合RabbitMQ 4.x消息队列到微服务体系,实现服务间异步通信、解耦,理解消息发送、接收、路由的完整流程。
  4. 培养微服务场景下的问题排查能力,能定位组件整合中的版本冲突、配置错误、通信异常等问题并解决。
  5. 建立微服务架构设计的初步思维,为后续复杂微服务系统开发、分布式事务处理等内容奠定基础。

三、实验完成内容概要(介绍实际完成的内容)

  1. 服务注册与发现:服务启动时注册到 Eureka,调用方从注册中心获取服务列表。
  2. 消息队列异步通信:生产者发送消息到 RabbitMQ,消费者监听队列异步处理,实现服务解耦、流量削峰。
  3. 版本绑定:Spring Cloud 2021.0.x 严格绑定 Spring Boot 2.7.x,保证依赖稳定。

四、实验过程记录(方法、操作截图与步骤说明)
application.yml

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;

    }

}
messageproducer

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;

    }

}
rabbitconfig

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

    }

}
usercontroller

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;

    }

}
步骤 2:启动 UserService(8081)

运行 userservice 启动类,激活 user profile。

查看 Eureka 控制台,确认 USERSERVICE 状态为 UP。


3.orderservice


package com.example.project5.microservice.orderservice;

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


package com.example.project5.microservice.orderservice;

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。


步骤 4:测试消息发送与接收

在浏览器访问: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>

实验六:微服务与信息中间件整合实践

一、实验目的和要求

  1. 掌握 Spring Cloud 与 RabbitMQ 4.2.3 的整合方法
  2. 理解消息中间件在微服务架构中的核心作用
  3. 熟悉 RabbitMQ 4.2.3 的核心特性在微服务中的应用
  4. 掌握分布式系统中消息可靠传递的实现方法
  5. 通过实践掌握 Spring Cloud Stream 与 RabbitMQ 的集成(适配新版函数式编程模型)

三、实验完成内容概要(介绍实际完成的内容)

  1. 搭建 Spring Boot + Spring Cloud Stream + RabbitMQ 项目。
  2. 配置并创建:普通队列、死信队列、惰性队列、延迟队列。
  3. 开发消息生产者(订单消息、延迟消息)。
  4. 开发消息消费者(支付消费者、死信消费者、延迟消费者)。
  5. 完成全链路验证,确保消息发送、存储、消费正常。
  6. 解决启动冲突、Bean 重复、符号找不到、版本兼容等问题。
  • 实验过程记录(方法、操作截图与步骤说明)
    config
    RabbitMQConfig

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

    }

}
DeadLetterConfig


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

    }

}
DelayedMessageConfig


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 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
PaymentConsumer

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

    }

}
DelayedMessageConsumer

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

    }

}
producer
OrderProducer

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;

    }

}
DelayedMessageProducer

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 秒";

    }

}
service
OrderService

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

        );

    }

}
PaymentService

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>
1.普通消息发送成功浏览器返回:

2.延迟消息发送成功浏览器返回:


控制台输出:

五、结论(结果与分析)

  1. 本实验成功实现了 Spring Cloud Stream + RabbitMQ 的整合,完成了普通消息、死信消息、惰性队列、延迟消息的全套功能。
  2. 消息生产者可通过 HTTP 接口正常发送消息,消费者可异步接收并处理消息,实现服务间解耦。
  3. 死信队列可正确接收被拒绝、过期或失败的消息,便于后续排查与重试,提升系统容错能力。
  4. 延迟队列基于死信机制实现稳定运行,可用于订单超时未支付、延迟通知等场景。
  5. 惰性队列将消息保存到磁盘,适合高并发、海量消息堆积场景,保护系统内存。
  6. 整体系统运行稳定,消息投递与消费可靠,满足微服务异步通信需求。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐