RocketMQ核心概念

发布时间:2026/8/24 1:49:06
RocketMQ核心概念 RocketMQ核心角色1. NameServer 名称服务器注册中心 路由中心无状态可集群部署。Broker 启动向所有 NameServer 注册Producer/Consumer 从 NameServer 获取 Broker、Topic 路由信息不存储消息只存元数据Topic、Broker、队列信息心跳检测Broker 定时上报NameServer 感知 Broker 下线。2. Broker 消息服务器Broker 内部组件CommitLog、ConsumeQueue、IndexFile。真正存储消息的节点RocketMQ 核心服务。接收 Producer 发送消息存储消息处理 Consumer 的拉取消息请求执行消息删除、重试、死信等逻辑。3. Producer 消息生产者发送消息的客户端业务系统向 Broker 发送消息。发送方式包括同步发送、异步发送、单向发送自动从 NameServer 获取路由负载均衡选择队列发送4. Consumer 消息消费者消费消息的客户端业务系统从 Broker 拉取消息消费。消费模式Pull主动拉 / Push封装后的 pull底层还是拉消费组ConsumerGroup同一组内多个实例共同消费一组消息。消息相关概念1. 主题Topic消息的分类消息的一级分类。一个 Topic 分布在多个 Broker每个 Broker 上会分配若干队列。2. 消息队列MessageQueueTopic 的物理分片类似 Kafka 的 Partition。一个 Topic 由多个 MessageQueue 组成发送消息时Producer 轮询 / 指定算法把消息分发到不同 Queue发送消息时Producer 轮询 / 指定算法把消息分发到不同 Queue3.Tag 标签消息二级过滤标签同一个 Topic 下细分。例Topic:order_topicTag:create /pay/cancel生产者设置 tag消费者可以subscribe(topic, pay||create)过滤服务端过滤减少网络传输。不满足标签的消息消费者不会接收但消息偏移量会正常移动。消息过滤除了通过tag标签还可以给消息加上用户自定义属性消费时通过MessageSelector类似sql方式实现灵活过滤。4. Message 消息体消息实体topic、tag、keys、body、properties、msgId。keys业务唯一键可用于查询消息msgId全局消息 ID。5. Offset 偏移量CommitLog Offset消息在 CommitLog 文件中的物理偏移全局唯一。Consumer Offset消费位点记录某个消费组在某个 MessageQueue 消费到哪条存储在 Broker。消费进度就是维护 offset提交 offset 之后代表消息消费完成。Broker 存储文件消息存储1. CommitLog真正存储全部消息内容的文件所有 topic 的消息全部顺序写入 CommitLog。文件固定大小默认 1G顺序写随机读高性能2. ConsumeQueue消费队列逻辑索引文件。每个 Topic 下每个 MessageQueue 对应一个 ConsumeQueue。只存commitLog offset 消息长度 tag hashcode。不存消息 body相当于 CommitLog 的索引。消费者消费时先读 ConsumeQueue 拿到物理地址再去 CommitLog 读取真正消息。3. IndexFile索引文件给 key/msgId 查询消息用可选。消费组 ConsumerGroup同一个消费组的所有消费者实例消费同一套消息1、集群消费Clustering【默认】Topic 下每条消息消费组内只消费一次。多个实例分摊队列做负载均衡。Topic 下每条消息消费组内只消费一次。多个实例分摊队列做负载均衡。队列数量是消费并发上限。Queue 数 ≥ 消费者实例数多余实例空闲。消费偏移量存储在broker已topicConsumerGroup为key,value(队列id:偏移量)为每个队列对应的消费偏移量。2、广播消费Broadcasting消费组下每一个消费者实例都会收到全部消息。不维护 offsetoffset 存在客户端。特殊消息1、重试队列 RETRY消费失败消息不会直接丢弃发往当前消费组对应的重试队列延迟重试。重试次数递增延迟超过最大重试次数进入死信队列。2、死信队列 DLQ多次重试消费仍然失败消息投递到死信队列人工处理。3、延时消息发送消息指定延时级别消息不会立刻投递到期后才真正进入业务队列。4、事务消息实现分布式事务半消息Half Message、提交、回查、回滚。关键机制1、负载均衡Producer发送消息轮询选择 MessageQueueConsumer客户端负载均衡分配 MessageQueue 给各个 consumer 实例。2、刷盘机制SYNC_FLUSH 同步刷盘消息落磁盘才返回可靠性高性能低ASYNC_FLUSH 异步刷盘写入 pagecache 就返回后台异步刷盘高性能生产常用3、复制机制主从SYNC_MASTER 同步双写ASYNC_MASTER 异步复制。消息的顺序消费消费者通过设置监听为MessageListenerOrderly可以保证同一个队列中的消息被顺序消费若是消息处理时报错消息不会被投递到重试队列而是等待一段时间后重新消费若是一直不成功消息消费会被堵塞造成消息堆积。但是不能完全保证所有消息都被顺序消费因为一个主题下可能会有多个队列。保证消息顺序消费首先需要设置监听为MessageListenerOrderly且保证一个主题下只有一个队列。若既要保证消息有序又要保证并发效率。例如一个订单流程包含创建-付款-推送-完成 每个步骤都有一个消息同一个订单需要保证有序消费,因为同一个订单orderId在每个步骤都是一样的。发送消息时可以通过MessageQueueSelector选择将同一个订单的消息发送到相同的队列从而保证一个订单的所有消息被顺序消费。批量消息生产者生产消息时可以批量发送需将消息封装到List同一批消息只能属于同一个topiclist中的数据不能超过4M这一批消息只能进入同一个topic的一个队列。/** * 批量消息-生产者 list不要超过4m */publicclassBatchProducer{publicstaticvoidmain(String[]args)throwsException{// 实例化消息生产者ProducerDefaultMQProducerproducernewDefaultMQProducer(BatchProducer);// 设置NameServer的地址producer.setNamesrvAddr(127.0.0.1:9876);// 启动Producer实例producer.start();StringtopicBatchTest;ListMessagemessagesnewArrayList();//1、这里必须是相同的topic否则会报错//2、这里list的数据不能超过4M否则会报错//3、这里list的数据只会进入topic的一个queue中for(inti0;i10000;i){messages.add(newMessage(topic,,null,(Hello world i).getBytes()));}try{producer.send(messages);}catch(Exceptione){producer.shutdown();e.printStackTrace();}// 如果不再发送消息关闭Producer实例。producer.shutdown();}}