前言:从基础到业务的闭环落地

经过前七篇的学习,我们已经掌握 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 补充日志发送场景(完善流程)

为了模拟完整的订单操作日志,在延时关单、库存扣减等逻辑中补充日志发送,确保同一订单日志顺序连贯:

  1. OrderConsumerServicereduceStock方法中添加日志:
private void reduceStock(String goodsId) {
    log.info("扣减库存成功,商品ID:{}", goodsId);
    // 发送扣减库存日志
    orderProducerService.sendOrderlyLogMsg(orderDTO.getOrderId(), "扣减库存", "商品库存扣减成功,商品ID:" + goodsId);
}
  1. DelayOrderConsumerServicecloseOrder方法中添加日志:
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 测试验证

  1. 调用下单接口/mq/order/create,触发订单创建、库存扣减、延时关单流程

  2. 查看控制台日志,确认同一订单的日志按“创建订单→扣减库存→取消订单”顺序消费

  3. 打开D:/rocketmq_logs/order_operation.log文件,验证日志按顺序写入,无乱序现象

四、三大场景核心避坑要点

通用避坑

  • 所有消息发送/消费逻辑,必须保证幂等性(如订单消息防重复消费、库存扣减防重复扣减),可通过订单ID、消息ID去重

  • 生产环境需关闭Broker自动创建Topic功能,提前规划Topic、Tag命名规范,避免乱建Topic导致运维混乱

  • 消息体尽量精简,避免传输大文件、大量冗余数据,提升发送/消费效率

分场景避坑

  • 订单异步处理:核心业务(订单入库)必须同步执行,非核心业务异步处理,消息发送采用同步发送,确保消息必达

  • 超时关单:延时等级需准确对应业务需求,测试用短延时,生产用30分钟(等级16),避免延时错误导致关单异常

  • 日志归集:顺序消费需开启ConsumeMode.ORDERLY,同一订单的日志必须用订单ID作为业务ID路由,避免乱序

五、入门篇完整总结

至此,RocketMQ入门篇8篇内容已全部完结。从基础认知到业务实战,我们完成了从“零基础”到“能落地”的完整学习,核心知识点总结如下:

5.1 核心知识体系

  1. 基础层:理解RocketMQ架构(NameServer、Broker、Producer、Consumer),掌握核心概念(Topic、Tag、Group、Queue、Offset)

  2. 环境层:完成Windows/Linux/Docker三种环境搭建,部署控制台,实现环境校验与问题排查

  3. 开发层:掌握原生Java Client、SpringBoot整合两种开发方式,实现普通消息、延时消息、顺序消息的收发

  4. 可靠性层:理解消费重试机制、死信队列原理,实现异常消息闭环处理,避免消息丢失、堆积

  5. 实战层:落地电商三大核心场景,串联所有知识点,掌握消息队列在实际业务中的应用思路

5.2 学习收获与后续方向

通过入门篇的学习,你已经具备RocketMQ基础开发与业务落地能力,能够解决中小型项目的消息队列需求。后续可向以下方向深入学习:

  • 进阶篇:RocketMQ集群部署、事务消息、消息过滤、消息回溯、监控告警等高级特性

  • 实战篇:分布式事务、削峰填谷、最终一致性等复杂场景落地,结合微服务架构整合

  • 运维篇:RocketMQ性能优化、问题排查、集群扩容、日志分析等运维技巧

消息队列是分布式系统的核心组件,RocketMQ作为国产高性能消息队列,在电商、金融、物流等领域应用广泛。入门只是起点,后续需结合实际业务多练、多排查,才能真正掌握其核心精髓。

感谢大家跟随RocketMQ入门篇的学习,祝各位在消息队列的学习与实践中,少踩坑、多成长!

Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐