SpringBoot 整合 RocketMQ 完整实战教程

RocketMQ 是阿里开源的分布式消息中间件,具备高吞吐、高可用、低延迟、支持事务消息、延时消息等特性,广泛应用于微服务异步通信、流量削峰、分布式事务解耦、日志收集等业务场景。本文将从零带你实现 SpringBoot 整合 RocketMQ,涵盖环境准备、依赖配置、生产者发送消息、消费者监听消息、常见消息类型实战及踩坑总结。

一、前置环境准备

整合前需提前搭建好 RocketMQ 服务,本地或服务器部署均可,核心需要两个组件:NameServer(注册中心)、Broker(消息存储与转发)。

1. 环境版本说明

  • SpringBoot:2.7.x(稳定通用版本)
  • RocketMQ:4.9.5(兼容主流SpringBoot版本)
  • JDK:1.8+
  • Maven:3.6+

2. 服务启动验证

启动 RocketMQ 的 NameServer 和 Broker 后,确保服务正常运行,默认端口:

  • NameServer:9876
  • Broker:10909、10911

可通过控制台命令验证服务状态,确保后续项目连接无异常。

二、引入Maven依赖

SpringBoot 整合 RocketMQ 核心依赖为 rocketmq-spring-boot-starter,无需手动配置原生客户端,starter 已自动封装核心配置。

<!-- RocketMQ SpringBoot 启动器 -->
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    &lt;version&gt;2.2.2&lt;/version&gt;
&lt;/dependency&gt;

版本适配说明:2.2.x 版本适配 SpringBoot2.x,若使用 SpringBoot3.x 需升级至 2.3.x 及以上版本。

三、全局配置文件配置

application.yml 中配置 RocketMQ 核心连接参数,包括 NameServer 地址、生产者组名,统一管理配置,便于环境切换。

# RocketMQ 配置
rocketmq:
  # NameServer 地址,集群环境用分号分隔
  name-server: 127.0.0.1:9876
  # 生产者配置
  producer:
    # 生产者组名,同一业务统一组名
    group: springboot-producer-group
    # 消息发送超时时间,默认3000ms
    send-timeout: 3000
    # 失败重试次数
    retry-times-when-send-failed: 2

四、消息生产者实现

RocketMQ 封装了 RocketMQTemplate 模板类,内置各类消息发送方法,可快速发送普通消息、同步消息、异步消息、对象消息。

1. 生产者工具类

import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

@Component
public class RocketMQProducer {

    // 注入RocketMQ模板类
    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    // 定义Topic主题,统一管理
    private static final String TEST_TOPIC = "springboot-test-topic";

    /**
     * 发送普通同步消息
     * 特点:发送后阻塞等待响应,可靠性高,适合重要业务场景
     */
    public void sendSyncMessage(String message) {
        rocketMQTemplate.syncSend(TEST_TOPIC, MessageBuilder.withPayload(message).build());
    }

    /**
     * 发送对象消息(自动序列化)
     */
    public <T> void sendObjectMessage(T data) {
        rocketMQTemplate.syncSend(TEST_TOPIC, MessageBuilder.withPayload(data).build());
    }

    /**
     * 发送延时消息
     * RocketMQ预设延时等级:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
     * level=3 代表10秒后消费
     */
    public void sendDelayMessage(String message) {
        rocketMQTemplate.syncSend(TEST_TOPIC, MessageBuilder.withPayload(message).build(), 3000, 3);
    }
}

2. 核心方法说明

  • syncSend:同步发送,阻塞等待消息投递结果,保证消息可靠送达,适用于支付、订单等核心业务
  • asyncSend:异步发送,不阻塞主线程,通过回调接收结果,适合高吞吐非核心业务
  • sendOneWay:单向发送,不等待响应、无回调,适用于日志收集、统计等无需确认的场景

五、消息消费者实现

消费者通过 @RocketMQMessageListener 注解监听指定 Topic,实现 RocketMQListener 接口处理消息,框架自动完成消息监听、消费、重试机制。

1. 普通消息消费者

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;

// 消费者监听注解:指定消费者组、监听的Topic
@RocketMQMessageListener(
        consumerGroup = "springboot-consumer-group",
        topic = "springboot-test-topic"
)
@Service
public class RocketMQConsumer implements RocketMQListener<String> {

    /**
     * 消息消费方法
     * 消息监听成功后自动执行,参数为生产者发送的消息内容
     */
    @Override
    public void onMessage(String message) {
        try {
            // 业务消费逻辑
            System.out.println("消费者成功接收消息:" + message);
        } catch (Exception e) {
            // 消费异常,RocketMQ会自动重试
            e.printStackTrace();
            throw new RuntimeException("消息消费失败");
        }
    }
}

2. 消费者核心注解参数说明

  • consumerGroup:消费者组名,同一组消费者负载均衡消费消息,不同组可重复消费同一消息
  • topic:监听的消息主题,必须与生产者发送的 Topic 一致
  • messageModel:消费模式,默认 CLUSTERING(集群模式,一条消息仅被组内一个消费者消费),可选 BROADCASTING(广播模式,组内所有消费者均消费)

六、接口测试验证

编写测试接口,调用生产者发送消息,验证整个收发流程是否正常。

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@RestController
@RequestMapping("/mq")
public class MqTestController {

    @Autowired
    private RocketMQProducer rocketMQProducer;

    @GetMapping("/send")
    public String sendMessage() {
        // 发送普通消息
        rocketMQProducer.sendSyncMessage("Hello SpringBoot RocketMQ 整合测试!");
        return "消息发送成功";
    }

    @GetMapping("/send/delay")
    public String sendDelayMessage() {
        rocketMQProducer.sendDelayMessage("这是一条10秒延时消息");
        return "延时消息发送成功";
    }
}

测试步骤

  1. 启动 SpringBoot 项目,确保无报错,消费者自动监听 Topic
  2. 访问接口 http://localhost:8080/mq/send
  3. 查看控制台,可打印出消费者接收的消息,证明整合成功
  4. 访问 /mq/send/delay,10秒后控制台打印延时消息,验证延时消息生效

七、对象消息收发实战

实际业务中常需要传递实体对象,RocketMQ starter 内置 Jackson 序列化,无需手动转换。

1. 定义实体类

import lombok.Data;
import java.io.Serializable;

@Data
public class User implements Serializable {
    private Long id;
    private String username;
    private String phone;
}

2. 对象消息消费者

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;

@RocketMQMessageListener(
        consumerGroup = "springboot-user-consumer-group",
        topic = "springboot-user-topic"
)
@Service
public class UserMessageConsumer implements RocketMQListener<User> {

    @Override
    public void onMessage(User user) {
        System.out.println("接收用户对象消息:" + user.toString());
    }
}

3. 新增对象消息发送接口

@GetMapping("/send/user")
public String sendUserMessage() {
    User user = new User();
    user.setId(1L);
    user.setUsername("测试用户");
    user.setPhone("13800138000");
    rocketMQProducer.sendObjectMessage(user);
    return "用户对象消息发送成功";
}

八、常见问题与踩坑总结

1. 连接失败:Cannot connect to NameServer

  • 原因:NameServer 地址配置错误、RocketMQ 服务未启动、防火墙端口未开放
  • 解决:核对 name-server 配置,开启9876、10911端口,重启服务

2. 消息重复消费

  • 原因:消费者消费后未正常返回、业务异常触发重试、集群模式下节点切换
  • 解决:业务实现幂等性(唯一ID去重、数据库幂等校验)

3. 消费者不监听消息

  • 原因:Topic 名称前后端不一致、消费者组名重复、注解未生效
  • 解决:统一生产者消费者 Topic,保证消费者类被 Spring 托管(添加@Service)

4. 序列化失败

  • 原因:实体类无空参构造、缺少序列化依赖
  • 解决:实体类添加无参构造,确保 lombok 依赖正常

九、总结

本文完成了 SpringBoot 与 RocketMQ 的完整整合,核心流程可概括为:引入starter依赖 → 配置NameServer和生产者参数 → 通过RocketMQTemplate发送消息 → 通过@RocketMQMessageListener监听消费消息

同时实现了普通消息、延时消息、对象消息三种常用消息类型,覆盖大部分业务场景。RocketMQ 凭借高可靠、高吞吐的特性,可完美解决微服务异步解耦、流量削峰、延时任务等核心需求,后续可在此基础上拓展事务消息、顺序消息、消息过滤等高级功能。


作 者:南烛
链 接:https://www.itnotes.top/archives/1430
来 源:IT笔记
文章版权归作者所有,转载请注明出处!


上一篇
下一篇