一、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 在以下场景表现优异:

  1. 高吞吐场景:日志收集、监控数据上报

  2. 分布式事务:金融交易、订单支付

  3. 实时计算:实时统计、风控系统

  4. 事件驱动架构:微服务解耦、数据同步

  5. 定时任务:延迟消息、任务调度

关键优势

  • 低延迟:99.6%的消息延迟在1ms内

  • 高可用:主从架构,自动故障转移

  • 高吞吐:单机支持10万级TPS

  • 可扩展:支持水平扩展

  • 功能完善:事务消息、顺序消息、定时消息

在实际应用中,需要根据业务特点选择合适的消息模式,并做好监控、告警和容错处理。

Logo

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

更多推荐