欢迎光临
我们一直在努力

中间件rabbitmq

RabbitMQ 完整实战教程(安装 + 原生使用 + SpringBoot+SpringCloud 微服务)

我会带你从零到一完整实现所有需求,步骤清晰、代码可直接运行,覆盖:

  • RabbitMQ 安装配置
  • 原生 Java 发送 / 接收消息
  • 7 种工作模式实战
  • SpringBoot 集成 RabbitMQ 实现业务
  • SpringCloud 微服务 + RabbitMQ 服务间通信

  • 一、RabbitMQ 安装与配置(Windows/Linux 通用)

    1. 前置依赖

    RabbitMQ 基于 Erlang 语言开发,必须先安装 Erlang。

    2. 下载安装

    • Erlang:https://www.erlang.org/downloads
    • RabbitMQ:https://www.rabbitmq.com/download.html

    3. 启动与开启管理控制台

    bash

    运行

    # 进入 RabbitMQ sbin 目录执行
    # 启动服务
    rabbitmq-server start

    # 开启可视化管理控制台(必开)
    rabbitmq-plugins enable rabbitmq_management

    4. 访问管理后台

    • 地址:http://localhost:15672
    • 默认账号 / 密码:guest / guest(仅限本地访问)

    5. 创建用户与权限(生产环境必备)

    bash

    运行

    # 创建用户 admin 密码 123456
    rabbitmqctl add_user admin 123456

    # 设置管理员权限
    rabbitmqctl set_user_tags admin administrator

    # 赋予所有虚拟机权限
    rabbitmqctl set_permissions -p / admin ".*" ".*" ".*"


    二、原生 Java 实现消息发送与接收(基础队列)

    1. 引入依赖

    xml

    <dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.18.0</version>
    </dependency>

    2. 消息生产者(发送)

    java

    运行

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;

    public class Producer {
    private final static String QUEUE_NAME = "simple_queue";

    public static void main(String[] argv) throws Exception {
    // 1. 创建连接工厂
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    factory.setUsername("admin");
    factory.setPassword("123456");

    // 2. 创建连接 和 通道
    try (Connection connection = factory.newConnection();
    Channel channel = connection.createChannel()) {

    // 3. 声明队列
    channel.queueDeclare(QUEUE_NAME, false, false, false, null);
    String message = "Hello RabbitMQ!";

    // 4. 发送消息
    channel.basicPublish("", QUEUE_NAME, null, message.getBytes());
    System.out.println("消息发送成功:" + message);
    }
    }
    }

    3. 消息消费者(接收)

    java

    运行

    import com.rabbitmq.client.*;

    public class Consumer {
    private final static String QUEUE_NAME = "simple_queue";

    public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    factory.setUsername("admin");
    factory.setPassword("123456");

    Connection connection = factory.newConnection();
    Channel channel = connection.createChannel();
    channel.queueDeclare(QUEUE_NAME, false, false, false, null);

    // 消费消息回调
    DeliverCallback deliverCallback = (consumerTag, delivery) -> {
    String message = new String(delivery.getBody(), "UTF-8");
    System.out.println("收到消息:" + message);
    };

    // 持续监听队列
    channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> {});
    }
    }


    三、RabbitMQ 7 种工作模式实战(核心)

    RabbitMQ 7 种模式:简单模式、工作队列、发布订阅、路由、主题、头部、RPC我给你最常用的 5 种核心实现(企业必备):

    1. 工作队列模式(Work Queues)

    场景:一个生产者,多个消费者,消息轮询分发

    • 生产者发送多条消息
    • 两个消费者平均分摊消费

    2. 发布订阅模式(Publish/Subscribe)

    场景:一条消息,所有消费者都能收到

    • 基于 扇形交换机(Fanout)
    • 交换机绑定队列,广播消息

    3. 路由模式(Routing)

    场景:根据路由键精准发送消息

    • 基于 直连交换机(Direct)
    • 例:error 日志只发给错误日志消费者

    4. 主题模式(Topic)

    场景:模糊匹配路由键(最常用)

    • 基于 主题交换机(Topic)
    • 通配符:# 匹配多个词,* 匹配一个词
    • 例:log.# 匹配所有日志消息

    5. RPC 模式

    场景:客户端发送消息 → 服务端处理 → 返回结果

    • 同步调用场景

    完整代码我已精简整合,你直接运行即可体验所有模式。


    四、SpringBoot 集成 RabbitMQ 实现业务逻辑

    1. 引入 SpringBoot AMQP 依赖

    xml

    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>

    2. application.yml 配置

    yaml

    spring:
    rabbitmq:
    host: localhost
    port: 5672
    username: admin
    password: 123456
    virtual-host: /

    3. 配置队列 + 交换机(主题模式为例)

    java

    运行

    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    public class RabbitConfig {
    // 队列
    public static final String BUSINESS_QUEUE = "business_queue";
    // 交换机
    public static final String BUSINESS_EXCHANGE = "business_exchange";
    // 路由键
    public static final String ROUTING_KEY = "business.#";

    // 声明队列
    @Bean
    public Queue businessQueue() {
    return new Queue(BUSINESS_QUEUE);
    }

    // 声明主题交换机
    @Bean
    public TopicExchange businessExchange() {
    return new TopicExchange(BUSINESS_EXCHANGE);
    }

    // 队列绑定交换机
    @Bean
    public Binding binding() {
    return BindingBuilder.bind(businessQueue())
    .to(businessExchange())
    .with(ROUTING_KEY);
    }
    }

    4. 业务消息生产者(Service)

    java

    运行

    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.stereotype.Service;
    import javax.annotation.Resource;

    @Service
    public class BusinessProducer {
    @Resource
    private RabbitTemplate rabbitTemplate;

    public void sendBusinessMsg(String msg) {
    // 发送消息:交换机 + 路由键 + 消息内容
    rabbitTemplate.convertAndSend(
    RabbitConfig.BUSINESS_EXCHANGE,
    "business.order",
    msg
    );
    System.out.println("业务消息发送:" + msg);
    }
    }

    5. 业务消息消费者(Service)

    java

    运行

    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Service;

    @Service
    public class BusinessConsumer {
    // 监听业务队列
    @RabbitListener(queues = RabbitConfig.BUSINESS_QUEUE)
    public void receive(String msg) {
    System.out.println("业务处理成功,收到消息:" + msg);
    // 在这里写你的业务逻辑:订单、支付、日志等
    }
    }

    6. 测试接口

    java

    运行

    import org.springframework.web.bind.annotation.GetMapping;
    import org.springframework.web.bind.annotation.RestController;
    import javax.annotation.Resource;

    @RestController
    public class TestController {
    @Resource
    private BusinessProducer producer;

    @GetMapping("/send")
    public String send() {
    producer.sendBusinessMsg("用户下单成功,订单号:" + System.currentTimeMillis());
    return "消息发送成功";
    }
    }


    五、SpringCloud 微服务 + RabbitMQ 实现服务间通信

    架构说明

    • 订单服务(生产者):创建订单 → 发送消息
    • 库存服务(消费者):监听消息 → 扣减库存
    • 基于 SpringCloud Alibaba / SpringCloud 通用

    1. 两个微服务都引入依赖

    xml

    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    <dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter</artifactId>
    </dependency>

    2. 公共配置(两个服务一致)

    yaml

    spring:
    rabbitmq:
    host: localhost
    username: admin
    password: 123456

    3. 订单服务(消息生产者)

    java

    运行

    @Service
    public class OrderService {
    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void createOrder() {
    // 1. 模拟创建订单
    String orderId = "ORDER_" + System.currentTimeMillis();
    System.out.println("订单创建成功:" + orderId);

    // 2. 发送消息给库存服务
    rabbitTemplate.convertAndSend("order_exchange", "order.stock", orderId);
    }
    }

    4. 库存服务(消息消费者)

    java

    运行

    @Service
    public class StockService {
    // 监听订单消息
    @RabbitListener(bindings = @QueueBinding(
    value = @Queue("stock_queue"),
    exchange = @Exchange("order_exchange"),
    key = "order.stock"
    ))
    public void reduceStock(String orderId) {
    System.out.println("库存服务扣减库存,订单号:" + orderId);
    // 库存业务逻辑
    }
    }

    5. 微服务调用效果

  • 启动订单服务、库存服务
  • 订单服务创建订单 → 发送 MQ 消息
  • 库存服务自动监听并消费消息 → 扣减库存
  • 微服务之间完全解耦

  • 六、完整运行总结

  • 启动 RabbitMQ,确保管理后台可访问
  • 运行原生 Java 代码:实现基础消息收发
  • 运行SpringBoot 项目:实现业务消息队列
  • 启动两个微服务:实现 SpringCloud + MQ 服务间通信
  • RabbitMQ 五大常用工作模式 完整可运行代码

    依赖统一用:

    xml

    <dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.18.0</version>
    </dependency>

    连接工具类(通用)

    java

    运行

    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;

    public class RabbitUtil {
    public static Connection getConnection() throws Exception{
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    factory.setPort(5672);
    factory.setUsername("admin");
    factory.setPassword("123456");
    return factory.newConnection();
    }
    }


    1. 简单模式(一对一)

    生产者

    java

    运行

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;

    public class SimpleProducer {
    public static final String QUEUE = "simple_queue";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    // 声明队列
    channel.queueDeclare(QUEUE,false,false,false,null);
    String msg = "简单模式消息";
    channel.basicPublish("",QUEUE,null,msg.getBytes());
    System.out.println("发送:"+msg);
    channel.close();
    conn.close();
    }
    }

    消费者

    java

    运行

    import com.rabbitmq.client.*;

    public class SimpleConsumer {
    public static final String QUEUE = "simple_queue";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.queueDeclare(QUEUE,false,false,false,null);
    DeliverCallback cb = (tag,msg)->{
    System.out.println("接收:"+new String(msg.getBody()));
    };
    channel.basicConsume(QUEUE,true,cb,tag->{});
    }
    }


    2. 工作队列模式 WorkQueues(一对多 轮询)

    生产者

    java

    运行

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;

    public class WorkProducer {
    public static final String QUEUE = "work_queue";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.queueDeclare(QUEUE,false,false,false,null);
    for (int i = 1; i <= 10; i++) {
    String msg = "任务"+i;
    channel.basicPublish("",QUEUE,null,msg.getBytes());
    System.out.println("发送任务:"+msg);
    }
    channel.close();
    conn.close();
    }
    }

    消费者 1

    java

    运行

    import com.rabbitmq.client.*;

    public class WorkConsumer1 {
    public static final String QUEUE = "work_queue";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.queueDeclare(QUEUE,false,false,false,null);
    DeliverCallback cb = (tag,msg)->{
    System.out.println("消费者1收到:"+new String(msg.getBody()));
    };
    channel.basicConsume(QUEUE,true,cb,tag->{});
    }
    }

    消费者 2

    java

    运行

    import com.rabbitmq.client.*;

    public class WorkConsumer2 {
    public static final String QUEUE = "work_queue";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.queueDeclare(QUEUE,false,false,false,null);
    DeliverCallback cb = (tag,msg)->{
    System.out.println("消费者2收到:"+new String(msg.getBody()));
    };
    channel.basicConsume(QUEUE,true,cb,tag->{});
    }
    }

    运行:先开两个消费者,再开生产者,消息均分


    3. 发布订阅模式 Fanout(广播 所有消费者都收)

    生产者

    java

    运行

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;

    public class FanoutProducer {
    public static final String EXCHANGE = "fanout_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    // 声明扇形交换机
    channel.exchangeDeclare(EXCHANGE,"fanout");
    String msg = "广播消息";
    channel.basicPublish(EXCHANGE,"",null,msg.getBytes());
    System.out.println("发送广播:"+msg);
    channel.close();
    conn.close();
    }
    }

    消费者 1

    java

    运行

    import com.rabbitmq.client.*;

    public class FanoutConsumer1 {
    public static final String EXCHANGE = "fanout_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"fanout");
    // 临时队列
    String queue = channel.queueDeclare().getQueue();
    // 绑定交换机
    channel.queueBind(queue,EXCHANGE,"");
    DeliverCallback cb = (tag,msg)->{
    System.out.println("消费者1收到广播:"+new String(msg.getBody()));
    };
    channel.basicConsume(queue,true,cb,tag->{});
    }
    }

    消费者 2

    java

    运行

    import com.rabbitmq.client.*;

    public class FanoutConsumer2 {
    public static final String EXCHANGE = "fanout_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"fanout");
    String queue = channel.queueDeclare().getQueue();
    channel.queueBind(queue,EXCHANGE,"");
    DeliverCallback cb = (tag,msg)->{
    System.out.println("消费者2收到广播:"+new String(msg.getBody()));
    };
    channel.basicConsume(queue,true,cb,tag->{});
    }
    }


    4. 路由模式 Routing Direct(精准匹配)

    生产者

    java

    运行

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;

    public class DirectProducer {
    public static final String EXCHANGE = "direct_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"direct");
    // 发不同路由键
    channel.basicPublish(EXCHANGE,"info",null,"普通日志".getBytes());
    channel.basicPublish(EXCHANGE,"error",null,"错误日志".getBytes());
    System.out.println("路由消息发送完成");
    channel.close();
    conn.close();
    }
    }

    消费者(只收 error)

    java

    运行

    import com.rabbitmq.client.*;

    public class DirectErrorConsumer {
    public static final String EXCHANGE = "direct_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"direct");
    String queue = channel.queueDeclare().getQueue();
    // 绑定error路由键
    channel.queueBind(queue,EXCHANGE,"error");
    DeliverCallback cb = (tag,msg)->{
    System.out.println("接收错误日志:"+new String(msg.getBody()));
    };
    channel.basicConsume(queue,true,cb,tag->{});
    }
    }

    消费者(收 info+error)

    java

    运行

    import com.rabbitmq.client.*;

    public class DirectAllConsumer {
    public static final String EXCHANGE = "direct_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"direct");
    String queue = channel.queueDeclare().getQueue();
    channel.queueBind(queue,EXCHANGE,"info");
    channel.queueBind(queue,EXCHANGE,"error");
    DeliverCallback cb = (tag,msg)->{
    System.out.println("接收全部日志:"+new String(msg.getBody()));
    };
    channel.basicConsume(queue,true,cb,tag->{});
    }
    }


    5. 主题模式 Topic(模糊匹配 最常用)

    通配符:

    • * 匹配一个单词
    • # 匹配多个单词

    生产者

    java

    运行

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;

    public class TopicProducer {
    public static final String EXCHANGE = "topic_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"topic");
    channel.basicPublish(EXCHANGE,"goods.add",null,"商品新增".getBytes());
    channel.basicPublish(EXCHANGE,"order.save",null,"订单保存".getBytes());
    System.out.println("主题消息发送完毕");
    channel.close();
    conn.close();
    }
    }

    消费者 1 匹配 goods.#

    java

    运行

    import com.rabbitmq.client.*;

    public class TopicGoodsConsumer {
    public static final String EXCHANGE = "topic_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"topic");
    String queue = channel.queueDeclare().getQueue();
    channel.queueBind(queue,EXCHANGE,"goods.#");
    DeliverCallback cb = (tag,msg)->{
    System.out.println("商品类消息:"+new String(msg.getBody()));
    };
    channel.basicConsume(queue,true,cb,tag->{});
    }
    }

    消费者 2 匹配 *.save

    java

    运行

    import com.rabbitmq.client.*;

    public class TopicSaveConsumer {
    public static final String EXCHANGE = "topic_ex";
    public static void main(String[] args) throws Exception {
    Connection conn = RabbitUtil.getConnection();
    Channel channel = conn.createChannel();
    channel.exchangeDeclare(EXCHANGE,"topic");
    String queue = channel.queueDeclare().getQueue();
    channel.queueBind(queue,EXCHANGE,"*.save");
    DeliverCallback cb = (tag,msg)->{
    System.out.println("保存类消息:"+new String(msg.getBody()));
    };
    channel.basicConsume(queue,true,cb,tag->{});
    }
    }


    运行顺序统一:

  • 启动 RabbitMQ 服务
  • 先启动所有消费者
  • 再启动生产者发送消息
  • RabbitMQ 全套常用依赖(Java/SpringBoot/SpringCloud 最全)

    一、原生 Java RabbitMQ 核心依赖

    xml

    <!– 原生rabbitmq客户端 –>
    <dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.21.0</version>
    </dependency>

    二、SpringBoot 整合 RabbitMQ(最常用)

    xml

    <!– SpringBoot AMQP 整合包 –>
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>

    <!– web测试接口必备 –>
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
    </dependency>

    <!– 实体类序列化、JSON –>
    <dependency>
    <groupId>com.alibaba</groupId>
    <artifactId>fastjson</artifactId>
    <version>1.2.83</version>
    </dependency>

    <!– lombok简化代码 –>
    <dependency>
    <groupId>org.projectlombok</groupId>
    <artifactId>lombok</artifactId>
    <optional>true</optional>
    </dependency>

    <!– 工具类 –>
    <dependency>
    <groupId>org.apache.commons</groupId>
    <artifactId>commons-lang3</artifactId>
    </dependency>

    三、SpringCloud 微服务 + RabbitMQ 全套依赖

    xml

    <!– 微服务注册中心 nacos(推荐) –>
    <dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-nacos-discovery</artifactId>
    </dependency>

    <!– 微服务配置中心 –>
    <dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-nacos-config</artifactId>
    </dependency>

    <!– 微服务远程调用 openfeign –>
    <dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-openfeign</artifactId>
    </dependency>

    <!– 网关gateway –>
    <dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-gateway</artifactId>
    </dependency>

    <!– 熔断降级 sentinel –>
    <dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-sentinel</artifactId>
    </dependency>

    <!– 消息队列 rabbitmq –>
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>

    四、MQ 高级功能必备依赖

    xml

    <!– 消息重试、死信队列、延迟队列必备 –>
    <dependency>
    <groupId>org.springframework.retry</groupId>
    <artifactId>spring-retry</artifactId>
    </dependency>

    <!– RabbitMQ 延迟队列插件依赖 –>
    <dependency>
    <groupId>com.github.rabbitmq</groupId>
    <artifactId>rabbitmq-delayed-message-exchange</artifactId>
    <version>3.10.0</version>
    </dependency>

    五、数据库 + 事务(业务实战必加)

    xml

    <!– mysql驱动 –>
    <dependency>
    <groupId>com.mysql</groupId>
    <artifactId>mysql-connector-j</artifactId>
    <scope>runtime</scope>
    </dependency>

    <!– mybatis-plus –>
    <dependency>
    <groupId>com.baomidou</groupId>
    <artifactId>mybatis-plus-boot-starter</artifactId>
    <version>3.5.3.1</version>
    </dependency>

    <!– 分布式事务 seata –>
    <dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-seata</artifactId>
    </dependency>

    六、测试依赖

    xml

    <!– springboot单元测试 –>
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-test</artifactId>
    <scope>test</scope>
    </dependency>

    七、父工程统一版本管理(直接复制)

    xml

    <modules>
    <module>eureka-server</module>
    <module>product-service</module>
    <module>order-service</module>
    <module>config-server</module>
    <module>gateway-service</module>

    </modules>
    <dependencyManagement>
    <dependencies>
    <!– SpringCloud Alibaba 版本 –>
    <dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-alibaba-dependencies</artifactId>
    <version>2021.0.5.0</version>
    <type>pom</type>
    <scope>import</scope>
    </dependency>
    <!– SpringCloud 版本 –>
    <dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-dependencies</artifactId>
    <version>2021.0.5</version>
    <type>pom</type>
    <scope>import</scope>
    </dependency>
    </dependencies>
    </dependencyManagement>

    RabbitMQ 完整 yml 配置(SpringBoot+SpringCloud 通用)

    基础最简配置

    yaml

    spring:
    # RabbitMQ配置
    rabbitmq:
    # 服务地址
    host: localhost
    # 默认端口
    port: 5672
    # 用户名
    username: admin
    # 密码
    password: 123456
    # 虚拟主机
    virtual-host: /

    完整版(消息确认、重试、手动签收、死信必备)

    yaml

    spring:
    rabbitmq:
    host: localhost
    port: 5672
    username: admin
    password: 123456
    virtual-host: /
    # 生产者配置
    publisher-confirm-type: correlated # 开启消息确认
    publisher-returns: true # 开启消息退回
    # 消费者配置
    listener:
    simple:
    acknowledge-mode: manual # 手动确认消息(推荐)
    prefetch: 1 # 每次只拿一条消息
    retry:
    enabled: true # 开启消费重试
    initial-interval: 1000 # 首次重试间隔1秒
    max-attempts: 3 # 最大重试3次
    multiplier: 2 # 重试间隔倍数

    常用参数说明

  • acknowledge-mode: manual手动签收,业务成功手动确认,失败拒绝,防止消息丢失
  • prefetch:1公平分发,谁处理快谁多拿任务
  • publisher-confirm-type生产者确认消息是否到达交换机
  • publisher-returns消息无法路由队列时退回生产者
  • 端口区分

    • 客户端通信端口:5672(代码连接用)
    • 后台管理页面端口:15672(浏览器访问)

    SpringCloud 微服务统一配置

    所有微服务MQ 配置完全一致,只改服务名即可

    yaml

    spring:
    application:
    name: order-service # 订单服务名
    rabbitmq:
    host: localhost
    port: 5672
    username: admin
    password: 123456
    virtual-host: /

    全局 pom.xml 版本锁定(直接用)

    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.15</version>
    <relativePath/>
    </parent>

    <groupId>com.rabbit</groupId>
    <artifactId>rabbit-demo</artifactId>
    <version>1.0.0</version>

    <properties>
    <maven.compiler.source>8</maven.compiler.source>
    <maven.compiler.target>8</maven.compiler.target>
    <spring-cloud-alibaba.version>2021.0.5.0</spring-cloud-alibaba.version>
    </properties>

    <dependencyManagement>
    <dependencies>
    <dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-alibaba-dependencies</artifactId>
    <version>${spring-cloud-alibaba.version}</version>
    <type>pom</type>
    <scope>import</scope>
    </dependency>
    </dependencies>
    </dependencyManagement>

    <dependencies>
    <!–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>
    <!–lombok–>
    <dependency>
    <groupId>org.projectlombok</groupId>
    <artifactId>lombok</artifactId>
    <optional>true</optional>
    </dependency>
    <!–测试–>
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-test</artifactId>
    <scope>test</scope>
    </dependency>
    </dependencies>
    </project>

    <!– 阿里云镜像仓库 –>

      <repositories>

        <repository>

          <id>aliyunmaven</id>

          <name>阿里云公共仓库</name>

          <url>https://maven.aliyun.com/repository/public</url>

          <releases>

            <enabled>true</enabled>

          </releases>

          <snapshots>

            <enabled>true</enabled>

          </snapshots>

        </repository>

      </repositories>

      <pluginRepositories>

        <pluginRepository>

          <id>aliyunmaven</id>

          <name>阿里云插件仓库</name>

          <url>https://maven.aliyun.com/repository/public</url>

          <releases>

            <enabled>true</enabled>

          </releases>

          <snapshots>

            <enabled>true</enabled>

          </snapshots>

        </pluginRepository>

      </pluginRepositories>

    <?xml version="1.0" encoding="UTF-8"?>

    <settings xmlns="http://maven.apache.org/SETTINGS/1.0.0"

              xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"

              xsi:schemaLocation="http://maven.apache.org/SETTINGS/1.0.0

              http://maven.apache.org/xsd/settings-1.0.0.xsd">

        <!– 阿里云镜像中央仓库,解决所有下载失败 –>

        <mirrors>

            <mirror>

                <id>aliyunmaven</id>

                <name>阿里云公共仓库</name>

                <url>https://maven.aliyun.com/repository/public</url>

                <mirrorOf>central</mirrorOf>

            </mirror>

            <mirror>

                <id>aliyunmaven2</id>

                <name>阿里云谷歌仓库</name>

                <url>https://maven.aliyun.com/repository/google</url>

                <mirrorOf>google</mirrorOf>

            </mirror>

            <mirror>

                <id>aliyunmaven3</id>

                <name>阿里云spring插件仓库</name>

                <url>https://maven.aliyun.com/repository/spring-plugin</url>

                <mirrorOf>spring-plugin</mirrorOf>

            </mirror>

        </mirrors>

        <profiles>

            <profile>

                <id>jdk-17</id>

                <activation>

                    <jdk>17</jdk>

                </activation>

                <repositories>

                    <repository>

                        <id>central</id>

                        <url>https://maven.aliyun.com/repository/public</url>

                        <releases>

                            <enabled>true</enabled>

                          </releases>

                        <snapshots>

                            <enabled>true</enabled>

                          </snapshots>

                    </repository>

                </repositories>

                <pluginRepositories>

                    <pluginRepository>

                        <id>central</id>

                        <url>https://maven.aliyun.com/repository/public</url>

                        <releases>

                            <enabled>true</enabled>

                          </releases>

                        <snapshots>

                            <enabled>true</enabled>

                          </snapshots>

                    </pluginRepository>

                </pluginRepositories>

            </profile>

        </profiles>

        <activeProfiles>

            <activeProfile>jdk-17</activeProfile>

        </activeProfiles>

    </settings>

    安装 RabbitMQ 延迟消息插件(Windows)

  • 下载rabbitmq_delayed_message_exchange-4.2.0.ez。
  • 复制到RabbitMQ Server\\rabbitmq_server-4.2.3\\plugins。
  • 管理员 CMD 执行:
  • bash

    运行

    rabbitmq-plugins enable rabbitmq_delayed_message_exchange
    rabbitmq-service stop
    rabbitmq-service start

  • 访问http://localhost:15672,确认x-delayed-message交换机类型存在。
  • package com.example.config;

    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    public class DeadLetterConfig {

    // 死信交换机(持久化)
    @Bean
    public DirectExchange deadLetterExchange() {
    return new DirectExchange("dead.letter.exchange", true, false);
    }

    // 死信队列(持久化,接收过期、拒收的消息)
    @Bean
    public Queue deadLetterQueue() {
    return QueueBuilder.durable("dead.letter.queue")
    .build();
    }

    // 正常队列(设置TTL,过期后进入死信队列)
    @Bean
    public Queue orderQueueWithTTL() {
    return QueueBuilder.durable("order.queue.ttl")
    .withArgument("x-message-ttl", 5000) // 消息过期时间:5秒
    .withArgument("x-dead-letter-exchange", "dead.letter.exchange") // 绑定死信交换机
    .withArgument("x-dead-letter-routing-key", "dead.letter.routing.key") // 死信路由键
    .build();
    }

    // 正常队列与订单交换机绑定
    @Bean
    public Binding orderBindingWithTTL() {
    return BindingBuilder.bind(orderQueueWithTTL()).to(orderExchange()).with("order.routing.key");
    }

    // 死信队列与死信交换机绑定
    @Bean
    public Binding deadLetterBinding() {
    return BindingBuilder.bind(deadLetterQueue()).to(deadLetterExchange()).with("dead.letter.routing.key");
    }

    // 复用订单交换机(避免重复创建)
    @Bean
    public DirectExchange orderExchange() {
    return new DirectExchange("order.exchange", true, false);
    }
    }

    package com.example.config;

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.CustomExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    import java.util.HashMap;
    import java.util.Map;

    @Configuration
    public class DelayedMessageConfig {

    // 延迟交换机(类型为 x-delayed-message,适配 RabbitMQ 4.2.3 延迟插件)
    @Bean
    public CustomExchange delayedExchange() {
    Map<String, Object> args = new HashMap<>();
    // 补充 RabbitMQ 4.2.3 必需参数:指定延迟交换机的底层类型(direct)
    args.put("x-delayed-type", "direct");
    // 参数:交换机名称、类型、持久化、自动删除、附加参数
    return new CustomExchange("delayed.exchange", "x-delayed-message", true, false, args);
    }

    // 延迟队列(持久化,避免重启丢失消息)
    @Bean
    public Queue delayedQueue() {
    return new Queue("delayed.queue", true);
    }

    // 绑定:延迟队列与延迟交换机绑定,指定路由键
    @Bean
    public Binding delayedBinding() {
    return BindingBuilder.bind(delayedQueue()).to(delayedExchange()).with("delayed.routing.key").noargs();
    }
    }

    package com.example.config;

    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    public class LazyQueueConfig {

    // 惰性队列配置(适合大量消息场景,消息优先存储到磁盘)
    @Bean
    public Queue lazyOrderQueue() {
    return QueueBuilder.durable("lazy.order.queue")
    .withArgument("x-queue-mode", "lazy") // 惰性队列标识
    .build();
    }

    @Bean
    public Queue lazyPaymentQueue() {
    return QueueBuilder.durable("lazy.payment.queue")
    .withArgument("x-queue-mode", "lazy")
    .build();
    }

    // 交换机配置
    @Bean
    public DirectExchange lazyOrderExchange() {
    return new DirectExchange("lazy.order.exchange", true, false);
    }

    @Bean
    public DirectExchange lazyPaymentExchange() {
    return new DirectExchange("lazy.payment.exchange", true, false);
    }

    // 绑定配置
    @Bean
    public Binding lazyOrderBinding() {
    return BindingBuilder.bind(lazyOrderQueue()).to(lazyOrderExchange()).with("lazy.order.routing.key");
    }

    @Bean
    public Binding lazyPaymentBinding() {
    return BindingBuilder.bind(lazyPaymentQueue()).to(lazyPaymentExchange()).with("lazy.payment.routing.key");
    }
    }

    package com.example.config;

    import org.springframework.amqp.core.*;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    public class RabbitMQConfig {

    // 声明交换机(强制覆盖)
    @Bean
    public DirectExchange orderExchange() {
    return ExchangeBuilder.directExchange("order.exchange")
    .durable(true)
    .build();
    }

    // 声明队列
    @Bean
    public Queue paymentQueue() {
    return QueueBuilder.durable("payment.queue").build();
    }

    // 绑定
    @Bean
    public Binding binding(Queue paymentQueue, DirectExchange orderExchange) {
    return BindingBuilder.bind(paymentQueue).to(orderExchange).with("order.key");
    }
    }

    package com.example.consumer;

    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    import java.util.function.Consumer;

    /**
    * 新版延迟消息消费者(函数式模型,替代旧版注解驱动)
    */
    @Configuration
    public class DelayedMessageConsumer {

    /**
    * 消费函数,方法名 delayedMessageInput 与 bootstrap.yml 中 function.definition 对应
    * 自动绑定 delayedMessageInput-in-0 通道,接收延迟消息
    */
    @Bean
    public Consumer<String> delayedMessageInput() {
    return delayedMessage -> {
    System.out.println("延迟消息已消费: " + delayedMessage);
    };
    }
    }

    package com.example.consumer;

    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    import java.util.function.Consumer;

    @Configuration
    public class PaymentConsumer {
    @Bean
    public Consumer<String> paymentInput() {
    return msg -> {
    System.out.println("支付服务已接收: " + msg);
    System.out.println("支付处理完成 ✅");
    };
    }
    }

    package com.example.controller;

    import com.example.service.OrderService;
    import org.springframework.http.MediaType;
    import org.springframework.web.bind.annotation.PostMapping;
    import org.springframework.web.bind.annotation.RequestBody;
    import org.springframework.web.bind.annotation.RestController;

    @RestController
    public class OrderController {

    private final OrderService orderService;

    public OrderController(OrderService orderService) {
    this.orderService = orderService;
    }

    @PostMapping(value = "/orders", consumes = MediaType.TEXT_PLAIN_VALUE)
    public String createOrder(@RequestBody String orderDetails) {
    orderService.createOrder(orderDetails);
    return "订单创建成功";
    }
    }

    package com.example.producer;

    import org.springframework.amqp.core.MessageDeliveryMode;
    import org.springframework.amqp.core.MessagePostProcessor;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.stereotype.Component;

    @Component
    public class DelayedMessageProducer {
    private final RabbitTemplate rabbitTemplate;

    public DelayedMessageProducer(RabbitTemplate rabbitTemplate) {
    this.rabbitTemplate = rabbitTemplate;
    }

    public void sendDelayedMessage(String message, int delayTime) {
    MessagePostProcessor processor = msg -> {
    msg.getMessageProperties().setHeader("x-delay", delayTime);
    // 正确写法:直接调用枚举,不需要写 DeliveryMode
    msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
    return msg;
    };

    rabbitTemplate.convertAndSend("delayed.exchange", "delayed.routing.key", message, processor);
    System.out.println("延迟消息已发送,延迟时间: " + delayTime + "ms,消息内容: " + message);
    }
    }

    package com.example.producer;

    import org.springframework.cloud.stream.function.StreamBridge;
    import org.springframework.messaging.support.MessageBuilder;
    import org.springframework.stereotype.Component;

    @Component
    public class OrderProducer {
    private final StreamBridge streamBridge;

    public OrderProducer(StreamBridge streamBridge) {
    this.streamBridge = streamBridge;
    }

    public void sendOrderMessage(String orderMessage) {
    boolean sent = streamBridge.send(
    "orderOutput-out-0",
    MessageBuilder.withPayload(orderMessage).build()
    );
    System.out.println("订单消息已发送: " + orderMessage);
    }
    }

    package com.example.service;

    import com.example.producer.OrderProducer;
    import org.springframework.stereotype.Service;

    @Service
    public class OrderService {

    // 构造器注入(替代 @Autowired,避免空指针风险)
    private final OrderProducer orderProducer;

    public OrderService(OrderProducer orderProducer) {
    this.orderProducer = orderProducer;
    }

    /**
    * 创建订单并发送消息
    * @param orderDetails 订单详情
    */
    public void createOrder(String orderDetails) {
    String orderMessage = "订单创建成功: " + orderDetails;
    orderProducer.sendOrderMessage(orderMessage);
    }
    }

    server:
    port: 8080

    spring:
    main:
    allow-bean-definition-overriding: true
    rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest

    # 强制 UTF-8 编码
    servlet:
    encoding:
    charset: UTF-8
    enabled: true
    force: true

    cloud:
    stream:
    # 强制指定函数,消除警告
    function:
    definition: paymentInput
    bindings:
    # 生产者
    orderOutput-out-0:
    destination: order.exchange
    contentType: text/plain
    # 消费者
    paymentInput-in-0:
    destination: order.exchange
    group: payment
    contentType: text/plain

    # 🔥 核心:强制声明交换机为 direct 类型,覆盖旧冲突
    rabbit:
    bindings:
    orderOutput-out-0:
    producer:
    exchange-type: direct
    routing-key: order.key
    paymentInput-in-0:
    consumer:
    binding-routing-key: order.key

    <?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>
    <groupId>com.example</groupId>
    <artifactId>microservice-rabbitmq-demo</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <name>microservice-rabbitmq-demo</name>
    <description>Demo project for RabbitMQ</description>

    <parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>3.3.2</version>
    <relativePath/>
    </parent>

    <properties>
    <java.version>17</java.version>
    <spring-cloud.version>2023.0.3</spring-cloud.version>
    </properties>

    <dependencies>
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    <dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-stream-binder-rabbit</artifactId>
    </dependency>
    <dependency>
    <groupId>org.projectlombok</groupId>
    <artifactId>lombok</artifactId>
    <optional>true</optional>
    </dependency>

    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-test</artifactId>
    <scope>test</scope>
    </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>
    <configuration>
    <excludes>
    <exclude>
    <groupId>org.projectlombok</groupId>
    <artifactId>lombok</artifactId>
    </exclude>
    </excludes>
    </configuration>
    </plugin>
    </plugins>
    </build>
    </project>

    赞(0)
    未经允许不得转载:171主机测评 » 中间件rabbitmq
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址