java

关注公众号 jb51net

关闭
首页 > 软件编程 > java > SpringBoot 整合 RocketMQ

SpringBoot 整合 RocketMQ 实现消息生产与消费的全过程

作者:希望永不加班

这篇文章给大家介绍SpringBoot整合RocketMQ实现消息生产与消费的全过程,本文结合实例代码给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧

一、为什么选 RocketMQ?

RocketMQ 是阿里开源的分布式消息中间件,经过双11高并发场景验证,核心优势是高吞吐量、高可靠性、低延迟,解决分布式系统中的异步解耦、流量削峰、数据一致性问题,比 RabbitMQ 更适合高并发业务场景。

1.1 核心作用

1.2 典型应用场景

1.3 核心组件

消息流转:生产者 → 主题(Topic) → 消费者,核心组件无需死记,理解作用即可:

二、环境准备

推荐 Docker 安装(一键启动,无需配置环境变量),无 Docker 环境可选择本地安装,步骤简洁可落地。

2.1 Docker 一键安装

前提:已安装 Docker,执行以下命令,一键启动 NameServer 和 Broker(整合在一起,简化操作):

# 拉取 RocketMQ 镜像(官方推荐版本,稳定兼容)
docker pull rocketmqinc/rocketmq:4.9.4
# 启动容器(映射端口,设置环境变量,简化配置)
docker run -d --name rocketmq -p 9876:9876 -p 10911:10911 -e "ROCKETMQ_NAMESRV_ADDR=localhost:9876" -e "ROCKETMQ_BROKER_ADDR=localhost:10911" rocketmqinc/rocketmq:4.9.4

验证:容器启动后,执行 docker ps,能看到 rocketmq 容器运行即成功(9876是 NameServer 端口,10911是 Broker 端口)。

2.2 本地安装

2.3 Maven 依赖

SpringBoot 提供 RocketMQ starter 依赖,无需引入复杂核心依赖,直接在 pom.xml 中添加以下依赖,一键整合(版本与 RocketMQ 对应,4.9.4 版本适配):

<?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.18</version>
        <relativePath/>
    </parent>
    <groupId>com.example</groupId>
    <artifactId>springboot-rocketmq-intro</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <properties>
        <java.version>1.8</java.version>
        <rocketmq.version>4.9.4</rocketmq.version>
    </properties>
    <dependencies>
        <!-- SpringBoot Web 依赖(用于编写接口测试生产者) -->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <!-- RocketMQ 整合依赖(核心) -->
        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-spring-boot-starter</artifactId>
            <version>${rocketmq.version}</version>
        </dependency>
        <!-- Lombok 简化代码(可选,与全系列风格统一) -->
        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
            <optional>true</optional>
        </dependency>
    </dependencies>
</project>

• 新建 cmd 窗口,进入 bin 文件夹,执行命令:start mqbroker.cmd -n localhost:9876 autoCreateTopicEnable=true

• 弹出窗口显示“Broker started successfully”即成功(autoCreateTopicEnable=true 表示自动创建主题,简化入门操作)。

三、SpringBoot 基础配置

在 application.yml 中配置 RocketMQ 连接信息(NameServer 地址、生产者组、消费者组),注释详细,入门阶段无需修改核心配置,直接复制即可:

spring:
  application:
    name: springboot-rocketmq-intro # 项目名称,可自定义
# RocketMQ 核心配置
rocketmq:
  name-server: localhost:9876 # NameServer 地址(本地默认,服务器填对应IP)
  producer:
    group: rocketmq-producer-group # 生产者组(自定义,同一组生产者共享配置)
    send-message-timeout: 3000 # 消息发送超时时间(毫秒)
    retry-times-when-send-failed: 1 # 发送失败重试次数(入门设1,避免重复发送)
  consumer:
    group: rocketmq-consumer-group # 消费者组(自定义,同一组消费者负载均衡消费)
    message-model: CLUSTERING # 消费模式:CLUSTERING(集群模式,默认)、BROADCASTING(广播模式)
    consume-thread-max: 10 # 最大消费线程数(入门无需修改)
    consume-timeout: 15000 # 消费超时时间(毫秒)

配置说明:

四、消息生产与消费

重点实现「普通消息」的生产与消费(入门首选),代码可直接复制,注释详细,跟着操作就能跑通,后续可扩展事务消息、延迟消息。

4.1 消息生产者(发送普通消息)

编写生产者,通过 RocketMQTemplate 发送消息到指定主题(Topic),这里用接口编写,方便通过浏览器测试,代码可直接复制:

import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RestController;
/**
 * RocketMQ 消息生产者
 * 发送普通消息到指定主题
 */
@RestController
public class RocketMQProducer {
    // 注入 RocketMQTemplate(SpringBoot 自动配置,直接使用)
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    // 自定义主题名称(可自定义,消费者需订阅相同主题)
    private static final String TOPIC_NAME = "springboot_rocketmq_topic";
    /**
     * 发送普通消息接口(可通过浏览器访问测试)
     * @param msg 要发送的消息内容(可自定义)
     * @return 发送结果
     */
    @GetMapping("/send/rocketmq/{msg}")
    public String sendMessage(@PathVariable String msg) {
        // 消息内容(可自定义,拼接前缀方便区分)
        String message = "【RocketMQ 普通消息】" + msg;
        // 发送消息:参数1=主题名称,参数2=消息内容
        rocketMQTemplate.convertAndSend(TOPIC_NAME, message);
        return "RocketMQ 消息发送成功!发送内容:" + message;
    }
}

4.2 消息消费者(订阅并消费消息)

编写消费者,通过 @RocketMQMessageListener 注解订阅指定主题和消费者组,当主题中有消息时,自动接收并处理,代码可直接复制:

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
/**
 * RocketMQ 消息消费者
 * 订阅指定主题,消费消息
 */
@Component // 交给Spring管理,否则无法订阅消息
@RocketMQMessageListener(
        topic = "springboot_rocketmq_topic", // 订阅的主题(必须与生产者的主题一致)
        consumerGroup = "${rocketmq.consumer.group}" // 消费者组(与配置文件中一致)
)
public class RocketMQConsumer implements RocketMQListener<String> {
    /**
     * 接收并处理消息(消息类型与生产者发送的一致,这里是String)
     * @param msg 接收的消息内容
     */
    @Override
    public void onMessage(String msg) {
        // 模拟消息处理(实际项目中替换为业务逻辑,如扣减库存、发送短信等)
        System.out.println("RocketMQ 消费者接收消息:" + msg);
        // 入门阶段暂不添加复杂业务逻辑,后续可扩展异常处理、消息确认等
    }
}

4.3 测试步骤

4.4 关键注意点

五、扩展:两种常用消息类型

入门掌握普通消息即可,以下两种常用消息类型可按需扩展,代码可直接复制到对应类中,无需修改核心配置。

5.1 同步消息(可靠发送,适合重要消息)

同步消息:生产者发送消息后,会等待 Broker 确认接收,确保消息发送成功(适合订单、支付等重要消息),添加到 RocketMQProducer 类中:

/**
 * 发送同步消息(可靠发送,等待Broker确认)
 */
@GetMapping("/send/rocketmq/sync/{msg}")
public String sendSyncMessage(@PathVariable String msg) {
    String message = "【RocketMQ 同步消息】" + msg;
    // 同步发送,返回发送结果(可判断是否发送成功)
    org.apache.rocketmq.client.producer.SendResult sendResult = rocketMQTemplate.syncSend(TOPIC_NAME, message);
    // 判断发送状态(SUCCESS 表示发送成功)
    if (sendResult.getSendStatus() == org.apache.rocketmq.client.producer.SendStatus.SEND_OK) {
        return "同步消息发送成功!发送内容:" + message + ",消息ID:" + sendResult.getMsgId();
    } else {
        return "同步消息发送失败!";
    }
}

5.2 异步消息(非阻塞,适合非重要消息)

异步消息:生产者发送消息后,无需等待确认,直接返回,通过回调函数处理发送结果(适合日志、通知等非重要消息),添加到 RocketMQProducer 类中:

/**
 * 发送异步消息(非阻塞,通过回调处理结果)
 */
@GetMapping("/send/rocketmq/async/{msg}")
public String sendAsyncMessage(@PathVariable String msg) {
    String message = "【RocketMQ 异步消息】" + msg;
    // 异步发送,回调函数处理发送结果
    rocketMQTemplate.asyncSend(TOPIC_NAME, message, new org.apache.rocketmq.spring.support.RocketMQSendCallback() {
        // 发送成功回调
        @Override
        public void onSuccess(org.apache.rocketmq.client.producer.SendResult sendResult) {
            System.out.println("异步消息发送成功,消息ID:" + sendResult.getMsgId());
        }
        // 发送失败回调
        @Override
        public void onException(Throwable e) {
            System.out.println("异步消息发送失败,原因:" + e.getMessage());
        }
    });
    return "异步消息发送中,请查看控制台回调结果!发送内容:" + message;
}

六、常见问题与解决方案

整理入门阶段最常遇到的问题,直接对照解决方案,避免卡壳,节省时间。

1. 项目启动报错:Could not connect to NameServer

• 原因1:RocketMQ 未启动(NameServer 或 Broker 未启动);

• 原因2:application.yml 中 name-server 配置错误(如端口错误、IP错误);

• 原因3:服务器安装的 RocketMQ,未开放 9876、10911 端口;

• 解决方案:启动 RocketMQ;核对 name-server 配置;开放服务器端口(阿里云/腾讯云安全组配置)。

2. 消息发送成功,但消费者收不到消息

• 原因1:生产者和消费者的主题(Topic)不一致;

• 原因2:消费者组配置错误(@RocketMQMessageListener 的 consumerGroup 与配置文件不一致);

• 原因3:消费者类未添加 @Co想在高并发场景下保证系统稳定?这篇文章手把手教你用SpringBoot整合RocketMQ实现消息生产与消费,从环境搭建到核心代码,覆盖异步解耦、流量削峰等关键应用,让你快速掌握可靠消息中间件并上手实战mponent 注解,Spring 未扫描到;

• 解决方案:核对主题名称;核对消费者组配置;给消费者类添加 @Component 注解,重启项目。

3. 启动 Broker 报错“内存不足”

• 原因:RocketMQ 默认 JVM 内存配置较大,本地环境内存不足;

• 解决方案:修改 bin 目录下的 runbroker.cmd(Windows)或 runbroker.sh(Linux),将 JVM 内存配置改小,如:
原配置:-Xms2g -Xmx2g -Xmn1g
修改后:-Xms512m -Xmx512m -Xmn256m,保存后重启 Broker。

4. 消息发送失败,提示“Topic does not exist”

学习本就是一个长期积累的过程,没有捷径,唯有坚持。希望这些干货能够真正帮到你,学以致用,不断提升,在自己的领域里越走越远。喜欢本文,别忘了点赞、在看、转发,我们下期干货继续!

• 原因:未开启自动创建主题,且主题未手动创建;

• 解决方案:启动 Broker 时添加参数 autoCreateTopicEnable=true(本地安装),或重新启动 Docker 容器(Docker 方式已配置)。

到此这篇关于SpringBoot 整合 RocketMQ 实现消息生产与消费的文章就介绍到这了,更多相关SpringBoot 整合 RocketMQ 内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多支持脚本之家!

您可能感兴趣的文章:
阅读全文