Skip to content

RocketMQ

RocketMQ 是 Apache 旗下的分布式消息中间件,常用于高吞吐异步通信、订单交易链路、日志收集、事件驱动架构、分布式事务最终一致性等场景。

它的核心特点可以概括为:

  • 高吞吐:适合大量消息写入和消费。
  • 可堆积:消费者短时间处理不过来时,消息可以先存储在 Broker 中。
  • 消费组:同一个主题可以被多个消费组独立消费。
  • 顺序消息:可以按业务键把消息路由到同一个队列,保证局部顺序。
  • 事务消息:通过半消息、二次确认、事务回查保证本地事务和消息发送最终一致。
  • 重试和死信:消费失败后可以自动重试,多次失败后进入死信队列。

整体架构

mermaid
flowchart TD
    P1["Producer"] --> NS["NameServer<br/>路由发现"]
    C1["Consumer"] --> NS
    P1 --> B1["Broker Master<br/>消息写入"]
    B1 --> B2["Broker Slave<br/>复制备份"]
    B1 --> C1
    NS -.->|路由注册与发现| B1

NameServer

NameServer 负责保存 Topic 路由信息。Producer 和 Consumer 启动后会连接 NameServer,查询某个 Topic 对应哪些 Broker、哪些队列。

NameServer 本身比较轻量,多个 NameServer 之间不做强一致同步,Broker 会定期向所有 NameServer 上报路由。

Broker

Broker 是真正存储和投递消息的服务端。它负责:

  • 接收 Producer 发来的消息。
  • 把消息写入 CommitLog。
  • 维护 Topic 下的队列信息。
  • 向 Consumer 投递消息。
  • 处理消费进度、重试、死信、事务回查等能力。

Producer

Producer 负责发送消息。常见发送方式有同步发送、异步发送、单向发送、顺序发送、事务发送。

业务上建议为消息设置清晰的 key,例如订单号、支付单号,方便排查和幂等。

Consumer

Consumer 负责消费消息。RocketMQ 中同一个 ConsumerGroup 内的多个消费者共同分担消息,不同 ConsumerGroup 之间相互独立。

消息发送到消费流程

mermaid
flowchart TD
    A[Producer 查询 Topic 路由] --> B[选择目标 MessageQueue]
    B --> C[发送消息到 Broker]
    C --> D[Broker 写入 CommitLog]
    D --> E[构建 ConsumeQueue 索引]
    E --> F[Consumer 拉取消息]
    F --> G[执行业务逻辑]
    G --> H{消费成功?}
    H -->|成功| I[提交消费进度]
    H -->|失败| J[进入重试队列]
    J --> F

学习路径

  1. 先看 领域模型,理解 Topic、MessageQueue、Producer、ConsumerGroup 的关系。
  2. 再看 订阅关系,理解消费组和订阅表达式为什么要保持一致。
  3. 需要保证同一业务键顺序时,看 顺序消息
  4. 需要主事务和消息最终一致时,看 事务消息
  5. 消费失败、幂等、死信处理,看 消费重试
  6. 如果正在从老版本迁移,看 新旧版本比较5.x 版本 API 使用

常见问题

参照官网文档本地安装时遇到的问题

  1. 在 RocketMQ 启动时遇到 交换内存问题
  2. 虚拟机环境可以参考 docker 安装 RocketMQ

为什么 RocketMQ 需要消费端幂等

RocketMQ 的消费语义更接近“至少消费一次”。只要出现消费超时、消费者宕机、ACK 失败、Broker 重新投递,就可能造成重复消费。

因此消费者不要依赖“消息只来一次”,而要用业务唯一键做幂等。例如订单支付成功事件可以用支付单号作为唯一键,写入数据库时利用唯一索引避免重复处理。

Topic 和 ConsumerGroup 怎么命名

命名要体现业务含义,不建议使用过于抽象的名称。

  • Topic 示例:ORDER_PAID_EVENTUSER_REGISTER_EVENT
  • ConsumerGroup 示例:POINTS_ORDER_PAID_CGSMS_ORDER_PAID_CG

Topic 表达“发生了什么事件”,ConsumerGroup 表达“哪个业务系统在消费这个事件”。

代码 Demo:同步发送普通消息

下面示例演示订单服务发送“订单已创建”事件。真实项目中要把 NameServer 地址、Topic、Group 放到配置文件中。

java
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;

import java.nio.charset.StandardCharsets;

public class RocketProducerDemo {
    public static void main(String[] args) throws Exception {
        DefaultMQProducer producer = new DefaultMQProducer("ORDER_PRODUCER_GROUP");
        producer.setNamesrvAddr("127.0.0.1:9876");
        producer.start();

        String body = "{\"orderId\":1001,\"userId\":2001}";
        Message message = new Message(
            "ORDER_EVENT",
            "ORDER_CREATED",
            "order-1001",
            body.getBytes(StandardCharsets.UTF_8)
        );

        SendResult result = producer.send(message);
        System.out.println(result);
        producer.shutdown();
    }
}

ORDER_EVENT 是 Topic,ORDER_CREATED 是 Tag,order-1001 是业务 Key。排查消息时,Key 很重要。