RocketMQ 消息队列
·
一、简介
RocketMQ 是一款高性能、高可用的分布式消息队列中间件。核心设计思想:将业务拆分为「主流程同步执行 + 次要流程异步执行」,实现业务解耦、流量削峰、异步处理。
二、核心应用场景
- 异步解耦:下单后发短信、消息推送、发放积分、记录日志,不阻塞主流程
- 流量削峰:大促、秒杀场景缓冲瞬时高流量,保护数据库不被压垮
- 数据同步:订单、用户数据同步至报表、财务、数据分析系统
- 延时任务:订单超时未支付自动关闭、业务超时提醒
- 分布式事务:配合消息实现跨服务最终一致性
三、基础架构流程
生产者(业务服务) → 投递消息到 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);
}
七、执行流程说明
- 用户发起下单请求,程序先同步保存订单,立刻返回「下单成功」给前端;
- 同时向 RocketMQ 投递一条消息;
- 消费者监听消息并异步执行发短信、加积分、记录日志等操作;
- 主流程不被阻塞,大幅提升接口响应速度。
八、技术总结
- 异步解耦:主业务与附属业务拆分,业务之间互不影响;
- 流量削峰:高并发场景下缓冲流量,降低数据库压力;
- 提升体验:前端无需等待所有业务执行完成,接口响应更快。
Redis分布式锁 + RocketMQ 组合实战(电商秒杀完整链路)
- 用户请求进入接口 → Redis分布式锁 控制并发,防止商品超卖;
- 扣减库存、生成订单(核心同步业务);
- 订单创建完成 → 发送 RocketMQ 消息;
- 消费者异步处理:短信通知、积分发放、日志记录。
技术分工
- Redis 分布式锁:解决分布式并发安全问题
- RocketMQ:解决异步解耦、高并发流量削峰问题
二者搭配是互联网高并发项目的经典技术组合。
更多推荐



所有评论(0)