一、简介

RocketMQ 是一款高性能、高可用的分布式消息队列中间件。核心设计思想:将业务拆分为「主流程同步执行 + 次要流程异步执行」,实现业务解耦、流量削峰、异步处理。

二、核心应用场景

  1. 异步解耦:下单后发短信、消息推送、发放积分、记录日志,不阻塞主流程
  2. 流量削峰:大促、秒杀场景缓冲瞬时高流量,保护数据库不被压垮
  3. 数据同步:订单、用户数据同步至报表、财务、数据分析系统
  4. 延时任务:订单超时未支付自动关闭、业务超时提醒
  5. 分布式事务:配合消息实现跨服务最终一致性

三、基础架构流程

生产者(业务服务) → 投递消息到 RocketMQ 服务端 → 消费者(监听服务) 异步消费执行业务

四、Maven 依赖

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.2.2</version>
</dependency>

五、配置文件 application.yml

rocketmq:
  name-server: 127.0.0.1:9876
  producer:
    group: order-producer-group

六、完整代码实战(订单异步通知案例)

1. 消息实体类

import lombok.Data;
import java.util.Date;

@Data
public class Order {
    private Long orderId;
    private Long userId;
    private String orderNo;
    private Date createTime;
}

2. 消息生产者(下单服务)

负责创建订单、发送异步消息,保证主流程快速响应。

import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;

@Service
public class OrderProducerService {

    @Resource
    private RocketMQTemplate rocketMQTemplate;

    @Resource
    private OrderMapper orderMapper;

    // 消息主题:生产者、消费者必须保持一致
    private static final String TOPIC = "order_notify_topic";

    /**
     * 创建订单并发送异步消息
     */
    public void createOrder(Order order) {
        // 1. 核心主业务:同步保存订单
        orderMapper.insertOrder(order);

        // 2. 发送异步消息,交由消费者处理附属业务
        rocketMQTemplate.syncSend(
                TOPIC,
                MessageBuilder.withPayload(order).build()
        );
        System.out.println("订单创建成功,消息发送完成");
    }
}

3. 消息消费者(异步处理服务)

监听消息主题,后台异步处理短信、积分、日志等非核心业务。

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

@Component
@RocketMQMessageListener(
        topic = "order_notify_topic",
        consumerGroup = "order_consumer_group"
)
public class OrderConsumer implements RocketMQListener<Order> {

    /**
     * 消息消费入口
     */
    @Override
    public void onMessage(Order order) {
        System.out.println("开始处理订单异步业务:" + order.getOrderNo());
        sendSms(order);
        addUserScore(order.getUserId());
        saveOrderLog(order);
    }

    // 模拟发送短信通知
    private void sendSms(Order order) {
        System.out.println("短信发送成功:" + order.getOrderNo());
    }

    // 模拟增加用户积分
    private void addUserScore(Long userId) {
        System.out.println("用户积分增加成功:" + userId);
    }

    // 模拟记录订单日志
    private void saveOrderLog(Order order) {
        System.out.println("订单日志记录完成");
    }
}

4. Order Mapper 数据层

import org.apache.ibatis.annotations.Insert;

public interface OrderMapper {
    @Insert("INSERT INTO `order`(order_id,user_id,order_no,create_time) " +
            "VALUES(#{orderId},#{userId},#{orderNo},#{createTime})")
    int insertOrder(Order order);
}

七、执行流程说明

  1. 用户发起下单请求,程序先同步保存订单,立刻返回「下单成功」给前端;
  2. 同时向 RocketMQ 投递一条消息;
  3. 消费者监听消息并异步执行发短信、加积分、记录日志等操作;
  4. 主流程不被阻塞,大幅提升接口响应速度。

八、技术总结

  1. 异步解耦:主业务与附属业务拆分,业务之间互不影响;
  2. 流量削峰:高并发场景下缓冲流量,降低数据库压力;
  3. 提升体验:前端无需等待所有业务执行完成,接口响应更快。

Redis分布式锁 + RocketMQ 组合实战(电商秒杀完整链路)

  1. 用户请求进入接口 → Redis分布式锁 控制并发,防止商品超卖;
  2. 扣减库存、生成订单(核心同步业务);
  3. 订单创建完成 → 发送 RocketMQ 消息
  4. 消费者异步处理:短信通知、积分发放、日志记录。

技术分工

  • Redis 分布式锁:解决分布式并发安全问题
  • RocketMQ:解决异步解耦、高并发流量削峰问题

二者搭配是互联网高并发项目的经典技术组合。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐