RocketMQ 原理与业务场景详解
一、RocketMQ 核心架构
1. 架构组件
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ Producer │ │ NameServer │ │ Consumer │
│ │ │ (注册中心) │ │ │
└────────┬────────┘ └────────┬────────┘ └────────┬────────┘
│ │ │
└──────────┐ ┌──────┘ ┌─────────────────┘
│ │ │
┌──────▼────▼───────────▼──────┐
│ Broker 集群 │
│ ┌─────────────────────┐ │
│ │ Master (主) │ │
│ │ ┌─────────────┐ │ │
│ │ │ CommitLog │ │ │
│ │ │ ConsumeQueue│ │ │
│ │ │ IndexFile │ │ │
│ └──┼─────────────┼────┘ │
│ │ Slave (从) │ │
│ └──────────────┘ │
└───────────────────────────────┘
2. 核心组件详解
NameServer
-
轻量级注册中心,负责 Broker 发现和路由管理
-
无状态设计,支持集群部署
-
数据存储在内存中,启动时从 Broker 获取路由信息
Broker
-
消息存储和转发服务器
-
主从架构:Master 处理读写,Slave 只读备份
-
存储核心文件:
-
CommitLog:顺序写,所有消息的物理存储
-
ConsumeQueue:消息逻辑队列,索引文件
-
IndexFile:消息索引,支持按 Key/时间查询
-
Producer
-
消息生产者,支持同步/异步/单向发送
-
内置负载均衡:轮询选择 MessageQueue
-
支持事务消息
Consumer
-
消息消费者,支持集群和广播模式
-
Pull 模型,支持长轮询
-
支持顺序消费和并发消费
二、核心工作原理
1. 消息存储机制
// CommitLog 顺序写入(高性能关键)
public class CommitLog {
// 所有消息按顺序写入,不区分 Topic
// 写入流程:
// 1. 获取写入锁
// 2. 追加到 MappedFile(内存映射文件)
// 3. 更新写入指针
// 4. 同步刷盘或异步刷盘
}
索引构建:
消息写入 CommitLog → 异步构建 ConsumeQueue → 消费者通过 ConsumeQueue 查找消息
2. 消息发送流程
// Producer 发送流程
public class DefaultMQProducer {
public SendResult send(Message msg) {
// 1. 校验消息
validateMessage(msg);
// 2. 查找路由(从 NameServer 获取 Topic 的路由信息)
TopicPublishInfo topicPublishInfo = tryToFindTopicPublishInfo(msg.getTopic());
// 3. 选择 MessageQueue(负载均衡)
MessageQueue mq = selectOneMessageQueue(topicPublishInfo, lastBrokerName);
// 4. 发送消息到 Broker
SendResult sendResult = sendKernelImpl(msg, mq, communicationMode, timeout);
return sendResult;
}
}
3. 消息消费流程
// PullConsumer 消费流程
public class DefaultMQPullConsumer {
public PullResult pull(MessageQueue mq, String subExpression,
long offset, int maxNums) {
// 1. 查找 Broker 地址
FindBrokerResult findBrokerResult = findBrokerResult(mq);
// 2. 拉取消息
PullResult pullResult = pullFromBroker(mq, subExpression,
offset, maxNums, findBrokerResult);
// 3. 更新消费进度
updateConsumeOffset(mq, pullResult.getNextBeginOffset());
return pullResult;
}
}
三、实际业务场景应用
场景1:电商订单系统(异步解耦)
业务需求:
-
用户下单后需要:扣减库存、生成订单、发送短信、增加积分
-
各服务之间需要解耦,提高系统可用性
实现方案:
// 1. 订单服务发送消息
@Component
public class OrderService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void createOrder(OrderDTO orderDTO) {
// 本地事务创建订单
Order order = orderMapper.create(orderDTO);
// 发送订单创建消息(事务消息)
Message<OrderCreatedEvent> message = MessageBuilder
.withPayload(new OrderCreatedEvent(order))
.setHeader(RocketMQHeaders.KEYS, order.getOrderNo())
.build();
rocketMQTemplate.sendMessageInTransaction(
"order-tx-group",
"ORDER_CREATED_TOPIC",
message,
order
);
}
}
// 2. 库存服务消费消息
@Component
@RocketMQMessageListener(
topic = "ORDER_CREATED_TOPIC",
consumerGroup = "inventory-consumer-group",
selectorExpression = "tag1" // 使用Tag过滤
)
public class InventoryConsumer implements RocketMQListener<OrderCreatedEvent> {
@Override
@Transactional
public void onMessage(OrderCreatedEvent event) {
// 扣减库存(保证幂等性)
inventoryService.deductStock(event.getOrder());
// 发送库存扣减成功消息
rocketMQTemplate.syncSend("INVENTORY_DEDUCTED_TOPIC",
new InventoryDeductedEvent(event.getOrder()));
}
}
场景2:秒杀系统(流量削峰)
业务需求:
-
瞬时高并发请求(如10万QPS)
-
避免数据库被压垮
-
保证公平性和一致性
实现方案:
@Component
public class SeckillService {
// 使用 RocketMQ 削峰填谷
public void handleSeckillRequest(SeckillRequest request) {
// 1. 快速验证(令牌桶限流)
if (!rateLimiter.tryAcquire()) {
throw new BusinessException("系统繁忙,请重试");
}
// 2. 发送消息到队列(异步处理)
Message<SeckillRequest> message = MessageBuilder
.withPayload(request)
.setHeader(RocketMQHeaders.KEYS, request.getSeckillId() + ":" + request.getUserId())
.build();
// 使用顺序消息,保证同一商品的请求顺序处理
rocketMQTemplate.syncSendOrderly(
"SECKILL_REQUEST_TOPIC",
message,
request.getSeckillId().toString()
);
}
}
// 秒杀处理器
@Component
@RocketMQMessageListener(
topic = "SECKILL_REQUEST_TOPIC",
consumerGroup = "seckill-processor-group",
consumeMode = ConsumeMode.ORDERLY // 顺序消费
)
public class SeckillProcessor implements RocketMQListener<SeckillRequest> {
@Override
public void onMessage(SeckillRequest request) {
// 3. 数据库扣减库存(控制并发度)
boolean success = seckillDao.reduceStock(request.getSeckillId());
if (success) {
// 4. 创建订单
orderService.createSeckillOrder(request);
// 5. 发送秒杀成功通知
rocketMQTemplate.asyncSend("SECKILL_SUCCESS_TOPIC",
new SeckillSuccessEvent(request));
}
}
}
场景3:数据同步(最终一致性)
业务需求:
-
主库数据变更同步到搜索引擎、缓存等
-
保证最终一致性
-
支持重试和补偿
实现方案:
// 使用 Canal + RocketMQ 实现数据同步
@Component
public class DataSyncService {
// 监听数据库变更
@EventListener
public void handleDataChange(DataChangeEvent event) {
// 构造变更消息
DataSyncMessage message = DataSyncMessage.builder()
.tableName(event.getTable())
.operation(event.getOperation())
.before(event.getBefore())
.after(event.getAfter())
.timestamp(System.currentTimeMillis())
.build();
// 发送延迟消息,确保事务提交后再同步
Message<DataSyncMessage> mqMessage = MessageBuilder
.withPayload(message)
.setHeader(RocketMQHeaders.KEYS, event.getKey())
.build();
// 延迟3秒,确保主事务提交
rocketMQTemplate.syncSend("DATA_SYNC_TOPIC",
mqMessage, 3000, 3); // 延迟级别3对应3秒
}
}
// ES同步消费者
@Component
@RocketMQMessageListener(
topic = "DATA_SYNC_TOPIC",
consumerGroup = "es-sync-consumer"
)
public class ESSyncConsumer implements RocketMQListener<DataSyncMessage> {
@Override
public void onMessage(DataSyncMessage message) {
try {
// 同步到Elasticsearch
switch (message.getOperation()) {
case INSERT:
case UPDATE:
esClient.index(message);
break;
case DELETE:
esClient.delete(message.getKey());
break;
}
} catch (Exception e) {
// 记录失败,进入重试队列
log.error("ES同步失败,消息进入重试队列", e);
throw e; // 抛出异常触发重试
}
}
}
场景4:分布式事务(事务消息)
业务需求:
-
跨服务的事务操作(如:订单+积分)
-
保证最终一致性
-
避免本地消息表方案
实现方案:
// 订单服务(事务消息生产者)
@Component
public class OrderTransactionService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Transactional
public void createOrderWithPoints(OrderDTO orderDTO) {
// 1. 创建订单(本地事务)
Order order = orderMapper.create(orderDTO);
// 2. 发送事务消息
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
"order-points-tx-group",
"ORDER_CREATED_FOR_POINTS",
MessageBuilder.withPayload(new PointsEvent(order))
.setHeader(RocketMQHeaders.KEYS, order.getOrderNo())
.build(),
order // 业务参数,用于回查
);
}
// 本地事务执行器
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
Order order = (Order) arg;
// 执行本地业务
orderService.finalizeOrder(order);
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK;
}
}
// 本地事务状态回查
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
String orderNo = msg.getHeaders().get(RocketMQHeaders.KEYS).toString();
Order order = orderMapper.selectByOrderNo(orderNo);
if (order.getStatus() == OrderStatus.COMPLETED) {
return RocketMQLocalTransactionState.COMMIT;
} else if (order.getStatus() == OrderStatus.CANCELLED) {
return RocketMQLocalTransactionState.ROLLBACK;
}
return RocketMQLocalTransactionState.UNKNOWN;
}
}
// 积分服务(消费者)
@Component
@RocketMQMessageListener(
topic = "ORDER_CREATED_FOR_POINTS",
consumerGroup = "points-service-consumer"
)
public class PointsConsumer implements RocketMQListener<PointsEvent> {
@Override
@Transactional
public void onMessage(PointsEvent event) {
// 增加用户积分(保证幂等)
pointsService.addPoints(event.getUserId(),
event.getOrderAmount(),
event.getOrderNo());
}
}
场景5:延迟任务
业务需求:
-
订单超时未支付自动取消
-
定时推送通知
-
会议开始前提醒
实现方案:
@Component
public class DelayTaskService {
// 发送延迟消息
public void scheduleOrderCancel(String orderNo, int delayMinutes) {
OrderCancelMessage message = new OrderCancelMessage(orderNo);
// RocketMQ 支持18个延迟级别:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
int delayLevel = convertMinutesToDelayLevel(delayMinutes);
rocketMQTemplate.syncSend("ORDER_CANCEL_DELAY_TOPIC",
MessageBuilder.withPayload(message)
.setHeader(RocketMQHeaders.KEYS, orderNo)
.build(),
delayLevel);
}
// 订单取消处理器
@Component
@RocketMQMessageListener(
topic = "ORDER_CANCEL_DELAY_TOPIC",
consumerGroup = "order-cancel-consumer"
)
public static class OrderCancelConsumer implements RocketMQListener<OrderCancelMessage> {
@Override
public void onMessage(OrderCancelMessage message) {
Order order = orderService.getOrder(message.getOrderNo());
// 检查订单状态,如果还是未支付则取消
if (order.getStatus() == OrderStatus.UNPAID) {
orderService.cancelOrder(message.getOrderNo(), "超时未支付");
// 释放库存
rocketMQTemplate.syncSend("INVENTORY_RELEASE_TOPIC",
new InventoryReleaseEvent(order));
}
}
}
}
四、最佳实践和优化
1. 生产环境配置
# 生产环境配置示例
# NameServer地址
rocketmq.name-server=192.168.1.100:9876;192.168.1.101:9876
# Producer配置
rocketmq.producer.group=PRODUCER_GROUP
rocketmq.producer.send-message-timeout=3000
rocketmq.producer.compress-message-body-threshold=4096
rocketmq.producer.retry-times-when-send-failed=2
rocketmq.producer.retry-times-when-send-async-failed=2
# Consumer配置
rocketmq.consumer.consume-thread-min=20
rocketmq.consumer.consume-thread-max=64
rocketmq.consumer.pull-batch-size=32
rocketmq.consumer.consume-message-batch-max-size=10
2. 监控和运维
# 查看集群状态
mqadmin clusterList -n localhost:9876
# 查看Topic信息
mqadmin topicStatus -t YOUR_TOPIC -n localhost:9876
# 查看消费者进度
mqadmin consumerProgress -g YOUR_CONSUMER_GROUP -n localhost:9876
# 发送测试消息
mqadmin sendMsg -t YOUR_TOPIC -p "test message" -n localhost:9876
3. 常见问题解决方案
消息丢失预防:
-
同步刷盘(Sync flush)
-
主从同步(SYNC_MASTER)
-
生产者确认机制
-
消费者手动ACK
消息堆积处理:
// 1. 增加消费者实例
// 2. 提高消费并行度
@RocketMQMessageListener(
topic = "BUSY_TOPIC",
consumerGroup = "consumer-group",
consumeThreadNumber = 32, // 增加消费线程
consumeMessageBatchMaxSize = 10 // 批量消费
)
// 3. 设置死信队列,转移积压消息
顺序消息保障:
// 发送顺序消息
rocketMQTemplate.syncSendOrderly("ORDER_TOPIC", message, "order123");
// 消费顺序消息
@RocketMQMessageListener(
consumeMode = ConsumeMode.ORDERLY, // 顺序消费
messageModel = MessageModel.CLUSTERING
)
五、总结
RocketMQ 在以下场景表现优异:
-
高吞吐场景:日志收集、监控数据上报
-
分布式事务:金融交易、订单支付
-
实时计算:实时统计、风控系统
-
事件驱动架构:微服务解耦、数据同步
-
定时任务:延迟消息、任务调度
关键优势:
-
低延迟:99.6%的消息延迟在1ms内
-
高可用:主从架构,自动故障转移
-
高吞吐:单机支持10万级TPS
-
可扩展:支持水平扩展
-
功能完善:事务消息、顺序消息、定时消息
在实际应用中,需要根据业务特点选择合适的消息模式,并做好监控、告警和容错处理。
更多推荐



所有评论(0)