RocketMQ 系列文章(入门篇第 8 篇):RocketMQ 电商核心场景实战
前言:从基础到业务的闭环落地
经过前七篇的学习,我们已经掌握 RocketMQ 的基础架构、环境搭建、核心概念、普通消息收发、SpringBoot 整合、延时 / 顺序消息,以及消费重试与死信队列**。本篇作为入门篇收官之作,聚焦电商领域三大核心场景 ——订单异步处理、超时关单自动触发、日志异步归集,将前文所有知识点串联落地,实现从 “会用 API” 到 “能解决实际业务问题” 的跨越。
本次实战基于 SpringBoot 环境,全程围绕真实电商业务逻辑编写代码,兼顾代码可读性与生产实用性,所有内容可直接复制到项目中运行验证。
前置准备:
完整 SpringBoot+RocketMQ 整合环境(沿用第 5-7 篇项目,无需重新搭建)
本地 RocketMQ 服务(NameServer+Broker)正常启动
提前创建 3 个核心 Topic:topic_order(订单主题)、topic_log(日志主题)、topic_delay_order(延时订单主题)
项目中引入 lombok、fastjson(JSON 序列化)依赖,简化开发
一、场景一:订单异步处理 —— 解耦核心业务与非核心逻辑
1.1 业务背景
电商下单流程中,创建订单是核心业务,而扣减库存、发送短信通知、记录操作日志属于非核心业务。若同步执行这些非核心逻辑,会大幅延长接口响应时间;通过 RocketMQ 异步处理,可将核心流程与非核心流程解耦,提升接口吞吐量,优化用户体验。
1.2 核心实现逻辑
生产者:接收前端下单请求,先完成订单入库(核心),再发送订单消息至topic_order,立即返回下单结果给用户
消费者:监听topic_order,异步执行扣减库存、发送短信(模拟)、更新订单状态等非核心逻辑
保障机制:消息发送采用同步发送(确保订单消息必达),消费逻辑保证幂等性(避免重复扣减库存)
1.3 代码实现
1.3.1 引入依赖(补充)
在pom.xml中添加 JSON 序列化依赖,用于消息体封装:
<!-- FastJSON JSON序列化/反序列化 -->
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.83</version>
</dependency>
1.3.2 定义订单实体类
创建OrderDTO,封装订单核心数据:
package com.rocketmq.entity;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.math.BigDecimal;
/**
* 订单数据传输对象
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class OrderDTO {
// 订单ID
private String orderId;
// 用户ID
private String userId;
// 商品ID
private String goodsId;
// 订单金额
private BigDecimal amount;
// 订单状态(0-待支付,1-已支付,2-已完成,3-已取消)
private Integer status;
}
1.3.3 订单生产者(核心业务 + 消息发送)
创建OrderProducerService,处理下单核心逻辑并发送消息:
package com.rocketmq.service;
import com.alibaba.fastjson.JSON;
import com.rocketmq.entity.OrderDTO;
import com.rocketmq.util.RocketMQUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.UUID;
@Slf4j
@Service
public class OrderProducerService {
// 订单主题(提前创建)
private static final String ORDER_TOPIC = "topic_order";
private static final String ORDER_TAG = "tag_order_create";
@Autowired
private RocketMQUtil rocketMQUtil;
/**
* 电商下单核心方法
* @param userId 用户ID
* @param goodsId 商品ID
* @param amount 订单金额
* @return 订单ID
*/
public String createOrder(String userId, String goodsId, BigDecimal amount) {
// 1. 核心业务:生成订单ID,模拟订单入库(实际项目中替换为数据库插入逻辑)
String orderId = UUID.randomUUID().toString().replace("-", "");
OrderDTO orderDTO = new OrderDTO(orderId, userId, goodsId, amount, 0);
log.info("核心业务完成,订单入库成功:{}", JSON.toJSONString(orderDTO));
// 2. 异步发送订单消息至RocketMQ,触发非核心业务处理
try {
rocketMQUtil.syncSend(ORDER_TOPIC, ORDER_TAG, JSON.toJSONString(orderDTO));
log.info("订单消息发送成功,orderId:{}", orderId);
} catch (Exception e) {
// 生产环境需增加异常告警(如短信/邮件通知),避免消息丢失
log.error("订单消息发送失败,orderId:{}", orderId, e);
throw new RuntimeException("订单创建失败,消息发送异常");
}
// 3. 返回订单ID给用户,无需等待非核心业务完成
return orderId;
}
}
1.3.4 订单消费者(异步处理非核心逻辑)
创建OrderConsumerService,监听topic_order,执行扣减库存、发送短信等逻辑:
package com.rocketmq.consumer;
import com.alibaba.fastjson.JSON;
import com.rocketmq.entity.OrderDTO;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
/**
* 订单消息消费者:处理异步非核心业务
*/
@Slf4j
@Component
@RocketMQMessageListener(
consumerGroup = "consumer_group_order",
topic = "topic_order",
selectorExpression = "tag_order_create",
maxReconsumeTimes = 3 // 最大重试3次,避免无限重试
)
public class OrderConsumerService implements RocketMQListener<String> {
/**
* 异步处理非核心业务:扣减库存、发送短信、记录日志
* @param message 订单消息(JSON格式)
*/
@Override
public void onMessage(String message) {
try {
// 1. 解析消息体
OrderDTO orderDTO = JSON.parseObject(message, OrderDTO.class);
log.info("开始消费订单消息:{}", JSON.toJSONString(orderDTO));
// 2. 非核心业务1:扣减库存(模拟,实际项目替换为数据库扣减逻辑)
reduceStock(orderDTO.getGoodsId());
// 3. 非核心业务2:发送短信通知(模拟,实际项目替换为第三方短信接口)
sendSms(orderDTO.getUserId(), orderDTO.getOrderId());
// 4. 非核心业务3:更新订单状态(模拟,实际项目替换为数据库更新逻辑)
log.info("非核心业务处理完成,orderId:{}", orderDTO.getOrderId());
} catch (Exception e) {
log.error("订单消息消费失败,消息:{}", message, e);
// 抛出异常触发重试,重试耗尽后进入死信队列
throw new RuntimeException("订单异步处理失败");
}
}
/**
* 模拟扣减库存
*/
private void reduceStock(String goodsId) {
log.info("扣减库存成功,商品ID:{}", goodsId);
}
/**
* 模拟发送短信通知
*/
private void sendSms(String userId, String orderId) {
log.info("短信发送成功,用户ID:{},订单ID:{}", userId, orderId);
}
}
1.3.5 测试接口
在MsgController中新增下单接口,验证场景一:
@Autowired
private OrderProducerService orderProducerService;
/**
* 电商下单测试接口
*/
@GetMapping("/order/create")
public String createOrderTest() {
// 模拟参数:用户ID=1001,商品ID=2001,金额=99.9
String orderId = orderProducerService.createOrder("1001", "2001", new BigDecimal("99.9"));
return "下单成功,订单ID:" + orderId + ",核心流程已完成,异步处理中";
}
二、场景二:超时关单自动触发 —— 延时消息落地
2.1 业务背景
电商中,用户下单后若30 分钟内未支付,需自动取消订单并回退库存,避免库存占用。RocketMQ 的延时消息是实现该场景的最佳方案,无需额外搭建定时任务,轻量高效。
2.2 核心实现逻辑
生产者:订单创建时,同步发送延时 30 分钟的消息至topic_delay_order
消费者:监听topic_delay_order,接收延时消息后,查询订单支付状态
业务判断:若订单仍为 “待支付”,执行取消订单、回退库存逻辑;若已支付,直接忽略
2.3 代码实现
2.3.1 延时消息生产者(扩展下单逻辑)
在OrderProducerService中新增延时关单消息发送方法:
/**
* 发送延时关单消息(30分钟后触发)
* @param orderDTO 订单信息
*/
public void sendDelayCloseOrderMsg(OrderDTO orderDTO) {
// 延时等级4=30秒(测试用,生产环境替换为等级16=30分钟)
// 注意:4.8.0版本延时等级最大为18,30分钟对应等级16,测试可先用等级4验证
int delayLevel = 4; // 测试用:30秒;生产用:16
try {
rocketMQUtil.sendDelayMsg(
"topic_delay_order",
"tag_delay_order_close",
JSON.toJSONString(orderDTO),
delayLevel
);
log.info("延时关单消息发送成功,orderId:{},延时等级:{}", orderDTO.getOrderId(), delayLevel);
} catch (Exception e) {
log.error("延时关单消息发送失败,orderId:{}", orderDTO.getOrderId(), e);
}
}
// 在createOrder方法末尾添加调用,实现“下单即触发延时关单”
// 2. 发送延时关单消息
sendDelayCloseOrderMsg(orderDTO);
2.3.2 延时关单消费者
创建DelayOrderConsumerService,监听延时消息并处理关单逻辑:
package com.rocketmq.consumer;
import com.alibaba.fastjson.JSON;
import com.rocketmq.entity.OrderDTO;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
/**
* 延时关单消费者:30分钟后自动检查订单支付状态
*/
@Slf4j
@Component
@RocketMQMessageListener(
consumerGroup = "consumer_group_delay_order",
topic = "topic_delay_order",
selectorExpression = "tag_delay_order_close",
maxReconsumeTimes = 2
)
public class DelayOrderConsumerService implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
try {
OrderDTO orderDTO = JSON.parseObject(message, OrderDTO.class);
log.info("收到延时关单消息,开始检查订单状态:{}", JSON.toJSONString(orderDTO));
// 1. 模拟查询订单支付状态(实际项目替换为数据库查询逻辑)
// 假设:模拟订单未支付,状态仍为0
if (orderDTO.getStatus() == 0) {
// 2. 执行关单逻辑:取消订单、回退库存
closeOrder(orderDTO.getOrderId());
rollbackStock(orderDTO.getGoodsId());
log.info("延时关单成功,orderId:{}", orderDTO.getOrderId());
} else {
log.info("订单已支付,无需关单,orderId:{}", orderDTO.getOrderId());
}
} catch (Exception e) {
log.error("延时关单消息消费失败,消息:{}", message, e);
throw new RuntimeException("延时关单处理失败");
}
}
/**
* 模拟取消订单
*/
private void closeOrder(String orderId) {
log.info("取消订单成功,orderId:{}", orderId);
}
/**
* 模拟回退库存
*/
private void rollbackStock(String goodsId) {
log.info("回退库存成功,商品ID:{}", goodsId);
}
}
2.3.3 测试验证
调用下单接口/mq/order/create,生成订单
观察控制台,30 秒后(测试用延时等级),消费者自动执行关单逻辑
若修改订单状态为 1(已支付),消费者会直接忽略关单逻辑,验证业务判断准确性
三、场景三:日志异步归集 —— 顺序消息 + 日志落地
3.1 业务背景
电商系统中,订单操作、支付操作、库存操作等日志需按时间顺序归集,便于后续审计、排查问题。若同步写入日志,会占用核心业务资源;通过顺序消息,可保证同一订单的操作日志按发送顺序消费,有序写入日志文件 / 数据库。
3.2 核心实现逻辑
生产者:同一订单的所有操作(创建、支付、关单),基于订单 ID发送顺序消息至topic_log
消费者:监听topic_log,采用顺序消费模式,保证同一订单的日志按顺序写入
落地效果:日志按业务顺序归集,避免日志乱序,提升审计效率
3.3 代码实现
3.3.1 定义日志实体类
创建LogDTO,封装日志核心数据:
package com.rocketmq.entity;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.time.LocalDateTime;
/**
* 操作日志DTO
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class LogDTO {
// 日志ID
private String logId;
// 关联订单ID
private String orderId;
// 操作类型(创建、支付、关单、扣减库存)
private String operationType;
// 操作时间
private LocalDateTime operationTime;
// 操作描述
private String description;
}
3.3.2 日志生产者(顺序消息发送)
在OrderProducerService中新增日志消息发送方法,基于订单 ID 实现顺序路由:
// 日志主题
private static final String LOG_TOPIC = "topic_log";
private static final String LOG_TAG = "tag_log_operation";
/**
* 发送顺序日志消息
* @param orderId 订单ID
* @param operationType 操作类型
* @param description 操作描述
*/
public void sendOrderlyLogMsg(String orderId, String operationType, String description) {
try {
// 基于订单ID发送顺序消息,保证同一订单的日志顺序一致
rocketMQUtil.sendOrderlyMsg(
LOG_TOPIC,
LOG_TAG,
JSON.toJSONString(new LogDTO(
UUID.randomUUID().toString().replace("-", ""),
orderId,
operationType,
LocalDateTime.now(),
description
)),
orderId // 业务ID:同一订单ID路由到同一Queue
);
log.info("顺序日志消息发送成功,orderId:{},操作类型:{}", orderId, operationType);
} catch (Exception e) {
log.error("顺序日志消息发送失败,orderId:{}", orderId, e);
}
}
// 在下单方法中新增日志发送,模拟创建订单日志
// 3. 发送创建订单日志消息
sendOrderlyLogMsg(orderId, "创建订单", "用户下单,创建订单成功");
3.3.3 日志消费者(顺序消费 + 日志落地)
创建LogConsumerService,监听topic_log,采用顺序消费模式,写入日志:
package com.rocketmq.consumer;
import com.alibaba.fastjson.JSON;
import com.rocketmq.entity.LogDTO;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import java.io.FileWriter;
import java.io.IOException;
import java.io.PrintWriter;
import java.time.format.DateTimeFormatter;
/**
* 日志消费者:顺序归集订单操作日志
*/
@Slf4j
@Component
@RocketMQMessageListener(
consumerGroup = "consumer_group_log",
topic = "topic_log",
selectorExpression = "tag_log_operation",
consumeMode = ConsumeMode.ORDERLY, // 开启顺序消费,保证同一订单日志有序
maxReconsumeTimes = 2
)
public class LogConsumerService implements RocketMQListener<String> {
// 日志文件路径(本地测试用,生产环境替换为日志框架路径)
private static final String LOG_FILE_PATH = "D:/rocketmq_logs/order_operation.log";
@Override
public void onMessage(String message) {
try {
// 1. 解析日志消息
LogDTO logDTO = JSON.parseObject(message, LogDTO.class);
log.info("收到顺序日志消息,准备写入文件:{}", JSON.toJSONString(logDTO));
// 2. 格式化日志内容(时间+订单ID+操作类型+描述)
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
String logContent = String.format(
"[%s] - 订单ID:%s,操作类型:%s,描述:%s%n",
logDTO.getOperationTime().format(formatter),
logDTO.getOrderId(),
logDTO.getOperationType(),
logDTO.getDescription()
);
// 3. 写入日志文件(追加模式,避免覆盖)
try (FileWriter fileWriter = new FileWriter(LOG_FILE_PATH, true);
PrintWriter printWriter = new PrintWriter(fileWriter)) {
printWriter.write(logContent);
} catch (IOException e) {
log.error("日志写入文件失败", e);
throw new RuntimeException("日志落地失败,触发重试");
}
log.info("日志顺序写入成功,日志内容:{}", logContent.trim());
} catch (Exception e) {
log.error("顺序日志消费失败,消息:{}", message, e);
throw new RuntimeException("日志处理失败,触发重试");
}
}
}
3.3.4 补充日志发送场景(完善流程)
为了模拟完整的订单操作日志,在延时关单、库存扣减等逻辑中补充日志发送,确保同一订单日志顺序连贯:
- 在
OrderConsumerService的reduceStock方法中添加日志:
private void reduceStock(String goodsId) {
log.info("扣减库存成功,商品ID:{}", goodsId);
// 发送扣减库存日志
orderProducerService.sendOrderlyLogMsg(orderDTO.getOrderId(), "扣减库存", "商品库存扣减成功,商品ID:" + goodsId);
}
- 在
DelayOrderConsumerService的closeOrder方法中添加日志:
private void closeOrder(String orderId) {
log.info("取消订单成功,orderId:{}", orderId);
// 发送取消订单日志
OrderProducerService orderProducerService = SpringContextUtil.getBean(OrderProducerService.class);
orderProducerService.sendOrderlyLogMsg(orderId, "取消订单", "订单超时未支付,自动取消");
}
注:SpringContextUtil为Spring上下文工具类,用于非Spring管理类获取Bean,生产环境可直接通过依赖注入获取OrderProducerService。
3.3.5 测试验证
-
调用下单接口
/mq/order/create,触发订单创建、库存扣减、延时关单流程 -
查看控制台日志,确认同一订单的日志按“创建订单→扣减库存→取消订单”顺序消费
-
打开
D:/rocketmq_logs/order_operation.log文件,验证日志按顺序写入,无乱序现象
四、三大场景核心避坑要点
通用避坑
-
所有消息发送/消费逻辑,必须保证幂等性(如订单消息防重复消费、库存扣减防重复扣减),可通过订单ID、消息ID去重
-
生产环境需关闭Broker自动创建Topic功能,提前规划Topic、Tag命名规范,避免乱建Topic导致运维混乱
-
消息体尽量精简,避免传输大文件、大量冗余数据,提升发送/消费效率
分场景避坑
-
订单异步处理:核心业务(订单入库)必须同步执行,非核心业务异步处理,消息发送采用同步发送,确保消息必达
-
超时关单:延时等级需准确对应业务需求,测试用短延时,生产用30分钟(等级16),避免延时错误导致关单异常
-
日志归集:顺序消费需开启
ConsumeMode.ORDERLY,同一订单的日志必须用订单ID作为业务ID路由,避免乱序
五、入门篇完整总结
至此,RocketMQ入门篇8篇内容已全部完结。从基础认知到业务实战,我们完成了从“零基础”到“能落地”的完整学习,核心知识点总结如下:
5.1 核心知识体系
-
基础层:理解RocketMQ架构(NameServer、Broker、Producer、Consumer),掌握核心概念(Topic、Tag、Group、Queue、Offset)
-
环境层:完成Windows/Linux/Docker三种环境搭建,部署控制台,实现环境校验与问题排查
-
开发层:掌握原生Java Client、SpringBoot整合两种开发方式,实现普通消息、延时消息、顺序消息的收发
-
可靠性层:理解消费重试机制、死信队列原理,实现异常消息闭环处理,避免消息丢失、堆积
-
实战层:落地电商三大核心场景,串联所有知识点,掌握消息队列在实际业务中的应用思路
5.2 学习收获与后续方向
通过入门篇的学习,你已经具备RocketMQ基础开发与业务落地能力,能够解决中小型项目的消息队列需求。后续可向以下方向深入学习:
-
进阶篇:RocketMQ集群部署、事务消息、消息过滤、消息回溯、监控告警等高级特性
-
实战篇:分布式事务、削峰填谷、最终一致性等复杂场景落地,结合微服务架构整合
-
运维篇:RocketMQ性能优化、问题排查、集群扩容、日志分析等运维技巧
消息队列是分布式系统的核心组件,RocketMQ作为国产高性能消息队列,在电商、金融、物流等领域应用广泛。入门只是起点,后续需结合实际业务多练、多排查,才能真正掌握其核心精髓。
感谢大家跟随RocketMQ入门篇的学习,祝各位在消息队列的学习与实践中,少踩坑、多成长!
更多推荐



所有评论(0)