Skip to content

RocketMQ

简介

RocketMQ 是阿里巴巴开源的分布式消息中间件,具有高可用、高性能、低延迟的特点。支持分布式事务消息、顺序消息、定时消息等高级特性。

核心概念

  • Producer:消息生产者,发送消息
  • Consumer:消息消费者,接收消息
  • Topic:主题,消息的第一级分类
  • Tag:标签,消息的第二级分类,用于精细化过滤
  • Message Queue:消息队列,Topic 的最小存储单元
  • Broker:消息存储和转发服务器

消息模型

Producer -> NameServer -> Broker -> Consumer
  • NameServer:无状态的路由发现中心,提供 Topic 路由信息
  • Broker:消息存储节点,主从架构保证高可用

消息类型

java
// 普通消息
DefaultMQProducer producer = new DefaultMQProducer("producer_group");
producer.start();
Message msg = new Message("TopicTest", "TagA", "Hello RocketMQ".getBytes());
SendResult result = producer.send(msg);

// 顺序消息
producer.send(msg, (mqs, msgObj, arg) -> {
    int index = (int) arg % mqs.size();
    return mqs.get(index);
}, orderId);

// 事务消息
TransactionMQProducer transactionProducer = new TransactionMQProducer("group");
transactionProducer.setTransactionListener(new TransactionListener() {
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        // 执行本地事务
        return LocalTransactionState.COMMIT_MESSAGE;
    }
    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        // 事务回查
        return LocalTransactionState.COMMIT_MESSAGE;
    }
});

适用场景

  • 异步解耦:订单系统与库存、物流系统解耦
  • 流量削峰:秒杀活动中的高并发流量缓冲
  • 日志处理:分布式系统的日志采集与分发
  • 分布式事务:基于 RocketMQ 的事务消息实现最终一致性