Administrator
发布于 2023-12-01 / 0 阅读
0
0

RocketMQ-202509120224

RocketMQ-202509120224

RocketMQ

1.RocketMQ简介

RocketMQ 是阿里巴巴开源的分布式消息中间件,现已成为 Apache 软件基金会的顶级项目。支持

事务消息、顺序消息、批量消息、定时消息、消息回溯等。

它里面有几个区别于标准消息中件间的概念,如Group、Topic、Queue等。

系统组成则由Producer、Consumer、Broker、NameServer等组件组成

2.RocketMQ 特点

是一个队列模型的消息中间件,具有高性能、高可靠、高实时、分布式等特点

Producer、Consumer、队列都可以分布式

Producer 向一些队列轮流发送消息,队列集合称为 Topic,Consumer 如果做广播消费,则一个

Consumer 实例消费这个 Topic 对应的所有队列,如果做集群消费,则多个 Consumer 实例平均

消费这个 Topic 对应的队列集合能够保证严格的消息顺序

支持拉(pull)和推(push)两种消息模式

pull其实就是消费者主动从MQ中去拉消息,而push则像rabbit MQ一样,是MQ给消费者推送消

息。

但是RocketMQ的push其实是基于pull来实现的。

它会先由一个业务代码从MQ中pull消息,然后再由业务代码push给特定的应用/消费者。其实底

层就是一个pull模式

高效的订阅者水平扩展能力

实时的消息订阅机制

亿级消息堆积能力

支持多种消息协议,如 JMS、OpenMessaging 等

较少的依赖

3.核心特性

高性能

1.单机万级队列:单台 Broker 节点支持创建超过一万个持久化队列,能够应对极高并发的生产

与消费场景。

2.顺序写磁盘:采用 Commit Log 结构,所有消息按序写入一个大文件,利用磁盘顺序写入的

高效率,实现高吞吐量和低延迟。

3.零拷贝技术:在消息读取过程中采用零拷贝技术,减少数据在内核空间和用户空间之间的复

制操作,提高 I/O 效率。

高可靠与高可用

1.分布式架构:系统由 NameServer、Broker、Producer 和 Consumer 四大组件构成,各组件

均可水平扩展,实现无单点故障。

2.多副本机制:支持主从副本模式,保证消息在 Broker 节点间的冗余备份,防止数据丢失。

3.故障自动切换:当主节点发生故障时,系统能够自动将流量切换至备节点,保证服务连续

性。

消息模型与丰富功能

1.多种消息类型:支持普通消息、顺序消息、事务消息、批量消息、定时(延时)消息、消息

回溯等,满足不同业务场景需求。

2.消息顺序性:提供单分区严格顺序消息和全局顺序消息,确保消息在特定场景下的消费顺

序。

3.事务消息:支持分布式事务,通过两阶段提交和回查机制,确保消息发送与数据库操作的最

终一致性。

易用性与灵活性

1.多种发送与消费模式:Producer 提供同步、异步、单向发送等多种方式;Consumer 支持拉

取(Pull)和推送(Push)两种消费模式。

2.丰富的客户端支持:提供 Java、Python、Go 等多种语言的 SDK,便于不同开发环境集成。

3.易于运维与管理:具备完善的监控指标、报警机制、命令行工具及图形化控制台,便于日常

运维和故障排查。

4.应用场景

RocketMQ 广泛应用于各种消息驱动的场景,包括但不限于:

订单系统:处理订单创建、支付、退款等流程中的异步通知和状态同步。

数据同步:在大数据处理中作为数据采集和分发的管道,实现数据的实时或准实时流转。

日志收集:作为日志聚合系统的一部分,接收、存储和转发各类应用产生的日志数据。

通知推送:用于发送短信、邮件、APP 推送等各类通知消息,实现用户触达。

任务调度:配合定时(延时)消息功能,执行定时任务或延迟任务。

5.RocketMQ 优势

目前主流的 MQ 主要是 RocketMQ、kafka、RabbitMQ,ActiveMQ其主要优势有:

1.支持事务型消息(消息发送和 DB 操作保持两方的最终一致性,RabbitMQ 和 Kafka 不支持)

支持结合 RocketMQ 的多个系统之间数据最终一致性(多方事务,二方事务是前提)

2.支持 18 个级别的延迟消息(Kafka 不支持)

3.支持指定次数和时间间隔的失败消息重发(Kafka 不支持,RabbitMQ 需要手动确认)

4.支持 Consumer 端 Tag 过滤,减少不必要的网络传输(即过滤由MQ完成,而不是由消费者

完成。RabbitMQ 和 Kafka 不支持)

5.支持重复消费(RabbitMQ 不支持,Kafka 支持)

6.在一个队列中可靠的先进先出(FIFO)和严格的顺序传递 (RocketMQ可以保证严格的消息

顺序,而ActiveMQ无法保证)

6.RocketMQ的原理

RocketMQ开发官方文档:

https://github.com/apache/rocketmq/blob/master/docs/cn/RocketMQ_Example.md

RocketMQ的集群架构如下

RocketMQ架构上主要分为四部分,如上图所示

6.1.Producer

消息发布的角色,支持分布式集群方式部署。Producer通过nameserver的负载均衡模块选择相应

的Broker集群队列进行消息投递,投递的过程支持快速失败并且低延迟。

6.2.Consumer

消息消费的角色,支持分布式集群方式部署。支持以push推,pull拉两种模式对消息进行消费。

同时 也支持集群方式和广播方式的消费,它提供实时消息订阅机制,可以满足大多数用户的需

求。

6.3.Broker

Broker主要负责消息的存储、投递和查询以及服务高可用保证。

6.4.NameServer

NameServer是一个Broker与Topic路由的注册中心支持Broker的动态注册与发现主要包括两个功

●Broker管理

NameServer接受Broker集群的注册信息并且保存下来作为路由信息的基本数据。然后提供心跳检

测机制,检查Broker是否还存活。

●路由信息管理

每个NameServer将保存关于Broker集群的整个路由信息和用于客户端查询的队列信息。然后Pro

ducer和Conumser通过NameServer就可以知道整个Broker集群的路由信息,从而进行消息的投

递和消费。

7.安装与配置

7.1下载地址

https://rocketmq.apache.org/download/

7.2解压压缩包,配置 ROCKETMQ_HOME

注意:rocketmq的配置默认不支持java8以上版本

在 mq的bin目录下找到并替换runbroker.cmd、runserver.cmd、tools.cmd三个启动脚本,如下:

runbroker.cmd

@echo off

rem Licensed to the Apache Software Foundation (ASF) under one or more

rem contributor license agreements. See the NOTICE file distributed

with

rem this work for additional information regarding copyright ownership.

rem The ASF licenses this file to You under the Apache License,

Version 2.0

rem (the "License"); you may not use this file except in compliance

with

rem the License. You may obtain a copy of the License at

rem

rem  http://www.apache.org/licenses/LICENSE-2.0

rem

runserver.cmd

rem Unless required by applicable law or agreed to in writing, software

rem distributed under the License is distributed on an "AS IS" BASIS,

rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or

implied.

rem See the License for the specific language governing permissions and

rem limitations under the License.

if not exist "%JAVA_HOME%\bin\java.exe" echo Please set the JAVA_HOME

variable in your environment, We need java(x64)! & EXIT /B 1

set "JAVA=%JAVA_HOME%\bin\java.exe"

setlocal

set BASE_DIR=%dp0

set BASE_DIR=%BASE_DIR:0,-1%

for %%d in (%BASE_DIR%) do set BASE_DIR=%%~dpd

set CLASSPATH=.;%BASE_DIR%lib*;%BASE_DIR%conf

rem

=======================================================================

====================

rem JVM Configuration

rem

=======================================================================

====================

set "JAVA_OPT=%JAVA_OPT% -server -Xms2g -Xmx2g"

rem set "JAVA_OPT=%JAVA_OPT% -XX:+UseG1GC -XX:G1HeapRegionSize=16m -

XX:G1ReservePercent=25 -XX:InitiatingHeapOccupancyPercent=30 -

XX:SoftRefLRUPolicyMSPerMB=0 -XX:SurvivorRatio=8"

set "JAVA_OPT=%JAVA_OPT% -XX:+UseZGC"

set "JAVA_OPT=%JAVA_OPT% -verbose:gc -Xlog:gc:%USERPROFILE%\mq_gc.log"

rem set "JAVA_OPT=%JAVA_OPT% -XX:+UseGCLogFileRotation -

XX:NumberOfGCLogFiles=5 -XX:GCLogFileSize=30m"

set "JAVA_OPT=%JAVA_OPT% -XX:-OmitStackTraceInFastThrow"

set "JAVA_OPT=%JAVA_OPT% -XX:+AlwaysPreTouch"

set "JAVA_OPT=%JAVA_OPT% -XX:MaxDirectMemorySize=15g"

rem set "JAVA_OPT=%JAVA_OPT% -XX:-UseLargePages -XX:-UseBiasedLocking"

set "JAVA_OPT=%JAVA_OPT% -XX:-UseLargePages"

rem set "JAVA_OPT=%JAVA_OPT% -

Djava.ext.dirs=%BASE_DIR%lib;%JAVA_HOME%\jre\lib\ext"

set "JAVA_OPT=%JAVA_OPT% -cp "%CLASSPATH%""

set "JAVA_OPT=%JAVA_OPT% --add-opens java.base/sun.nio.ch=ALL-UNNAMED"

"%JAVA%" %JAVA_OPT% %*

@echo off

rem Licensed to the Apache Software Foundation (ASF) under one or more

tools.cmd

rem contributor license agreements. See the NOTICE file distributed

with

rem this work for additional information regarding copyright ownership.

rem The ASF licenses this file to You under the Apache License,

Version 2.0

rem (the "License"); you may not use this file except in compliance

with

rem the License. You may obtain a copy of the License at

rem

rem  http://www.apache.org/licenses/LICENSE-2.0

rem

rem Unless required by applicable law or agreed to in writing, software

rem distributed under the License is distributed on an "AS IS" BASIS,

rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or

implied.

rem See the License for the specific language governing permissions and

rem limitations under the License.

if not exist "%JAVA_HOME%\bin\java.exe" echo Please set the JAVA_HOME

variable in your environment, We need java(x64)! & EXIT /B 1

set "JAVA=%JAVA_HOME%\bin\java.exe"

setlocal

set BASE_DIR=%dp0

set BASE_DIR=%BASE_DIR:0,-1%

for %%d in (%BASE_DIR%) do set BASE_DIR=%%~dpd

set CLASSPATH=.;%BASE_DIR%lib*;%BASE_DIR%conf

set "JAVA_OPT=%JAVA_OPT% -server -Xms2g -Xmx2g -Xmn1g -

XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m"

rem set "JAVA_OPT=%JAVA_OPT% -XX:+UseConcMarkSweepGC -

XX:+UseCMSCompactAtFullCollection -

XX:CMSInitiatingOccupancyFraction=70 -XX:+CMSParallelRemarkEnabled -

XX:SoftRefLRUPolicyMSPerMB=0 -XX:+CMSClassUnloadingEnabled -

XX:SurvivorRatio=8 -XX:-UseParNewGC"

set "JAVA_OPT=%JAVA_OPT% -XX:+UseZGC"

set "JAVA_OPT=%JAVA_OPT% -verbose:gc -

Xlog:gc:"%USERPROFILE%\rmq_srv_gc.log""

set "JAVA_OPT=%JAVA_OPT% -XX:-OmitStackTraceInFastThrow"

set "JAVA_OPT=%JAVA_OPT% -XX:-UseLargePages"

rem set "JAVA_OPT=%JAVA_OPT% -

Djava.ext.dirs=%BASE_DIR%lib;%JAVA_HOME%\jre\lib\ext"

set "JAVA_OPT=%JAVA_OPT% -cp "%CLASSPATH%""

" JAVA "

JAVA OPT

*

@echo off

7.3.启动MQ

rem Licensed to the Apache Software Foundation (ASF) under one or more

rem contributor license agreements. See the NOTICE file distributed

with

rem this work for additional information regarding copyright ownership.

rem The ASF licenses this file to You under the Apache License,

Version 2.0

rem (the "License"); you may not use this file except in compliance

with

rem the License. You may obtain a copy of the License at

rem

rem  http://www.apache.org/licenses/LICENSE-2.0

rem

rem Unless required by applicable law or agreed to in writing, software

rem distributed under the License is distributed on an "AS IS" BASIS,

rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or

implied.

rem See the License for the specific language governing permissions and

rem limitations under the License.

if not exist "%JAVA_HOME%\bin\java.exe" echo Please set the JAVA_HOME

variable in your environment, We need java(x64)! & EXIT /B 1

set "JAVA=%JAVA_HOME%\bin\java.exe"

setlocal

set BASE_DIR=%dp0

set BASE_DIR=%BASE_DIR:0,-1%

for %%d in (%BASE_DIR%) do set BASE_DIR=%%~dpd

set CLASSPATH=.;%BASE_DIR%lib*;%BASE_DIR%conf

rem

=======================================================================

====================

rem JVM Configuration

rem

=======================================================================

====================

set "JAVA_OPT=%JAVA_OPT% -server -Xms1g -Xmx1g -Xmn256m -

XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=128m"

rem set "JAVA_OPT=%JAVA_OPT% -

Djava.ext.dirs="%BASE_DIR%\lib";"%JAVA_HOME%\jre\lib\ext";"%JAVA_HOME%\

lib\ext""

set "JAVA_OPT=%JAVA_OPT% -cp "%CLASSPATH%""

"%JAVA%" %JAVA OPT% %*

启动NameServer

Cmd命令框执行进入至‘MQ文件夹\bin’下,然后执行 start mqnamesrv.cmd,启动NameServer。

启动成功请勿关闭此窗口!

启动Broker

进入至‘MQ文件夹\bin’下,CMD执行start mqbroker.cmd -n 127.0.0.1:9876 autoCreateTopicEna

ble=true ,启动Broker。启动成功请勿关闭此窗口!

8.RocketMQ入门

官方案例:链接

8.1.导入依赖

注意和安装的MQ版本一致

8.2.生产者

org.apache.rocketmq

rocketmq-client

5.0.0

步骤分析

1.创建producer组

2.设置NameServer地址

3.startr生产者

4.发送消息获取结果

5.结束producer

import org.apache.rocketmq.client.producer.MessageQueueSelector;

import org.apache.rocketmq.client.producer.SendCallback;

import org.apache.rocketmq.common.message.Message;

import org.apache.rocketmq.client.producer.DefaultMQProducer;

import org.apache.rocketmq.client.producer.SendResult;

import org.apache.rocketmq.common.message.MessageQueue;

import java.nio.charset.StandardCharsets;

import java.util.List;

//消息发送者

public class ProducerTest {

public static void main(String[] args) {

// 实例化消息生产者Producer MQ生产者 , 可以指定组名

producerGroupName

DefaultMQProducer producer = new

DefaultMQProducer("producergroup7");

// 设置NameServer的地址

producer.setNamesrvAddr("localhost:9876");

// 设置发送超时时间(毫秒)

producer.setSendMsgTimeout(5000);

// 设置重试次数

producer.setRetryTimesWhenSendAsyncFailed(3);

/**

  • 设置队列数量为1,默认为4 初始创建设置为1,消费者设置最大线程为1,实现
  • 顺序消费

  • 当 Producer发送消息或 Consumer 订阅一个 不存在的 Topic 时,
  • RocketMQ 会 自动创建该 Topic,并使用 defaultTopicQueueNums 作为队列数。

  • 如果 Topic 已经存在,此设置不会生效。
  • */

    producer.setDefaultTopicQueueNums(2); try {

    // 启动Producer实例

    producer.start();

    String[] tags = new String[]{"tags_error"};

    int sendCount = 0;

    for (int i = 0; i < 20; i++) {

    /**

  • 构建消息 ,参数为:topic,tags,内容
  • Topic即主题,是消息队列中用于对消息进行分类和组织的一种机制
  • Tag即标签,是消息队列中用于标记消息的一种属性或标识。它可以理
  • 解为Topic的子集,可以设置多个,用空格隔开

    */

    Message message = null;

    for (String tag : tags) {

    message = new Message("topic_log7", tag, ("我是消

    息" + i).getBytes());

    /**

  • SendResult :发送结果,其中包含
  • sendStatus=SEND_OK :发送状态
  • msgId :producer 创建的消息ID
  • offsetMsgId :Brocker创建的消息ID
  • messageQueue :消息存储的队列
  • */

    // 同步发送消息

    // synchronizationMessage(producer, message);

    // 异步发送消息

    // asynchronousMessage(producer,message);

    // 同样id数据放同一队列发送消息

    //for (int j = 0; j < 2; j++) {

    // message = new Message("topic_log", tag, ("我是

    消息:" + i + "-" + j).getBytes());

    // orderMessage(producer, message, i);

    //}

    // 延时消息

    delayMessage(producer, message);

    sendCount++;

    }

    }

    System.out.println("发送成功次数:" + sendCount);

    } catch (Exception e) {

    e.printStackTrace();

    } finally {

    // 如果不再发送消息,关闭Producer实例。

    producer.shutdown();

    }

    }

    /**

  • 同步发送消息
  • */

    public static void synchronizationMessage(DefaultMQProducer

    producer, Message message) throws Exception {

    SendResult send = producer.send(message);

    //System.out.printf("同步发送成功:%s%n", send);

    System.out.printf("同步发送成功,queueId:%d,topic:%s,tags:%s%n",

    send.getMessageQueue().getQueueId(),

    message.getTopic(), message.getTags());

    }

    /**

  • 异步发送消息
  • */

    public static void asynchronousMessage(DefaultMQProducer producer,

    Message message) throws Exception {

    producer.send(message, new SendCallback() {

    @Override

    public void onSuccess(SendResult sendResult) {

    System.out.printf("异步发送成功:%s%n", sendResult);

    }

    @Override

    public void onException(Throwable throwable) {

    System.out.println("异步发送失败");

    }

    });

    }

    /**

  • 单向发送
  • */

    public static void oneWayMessage(DefaultMQProducer producer,

    Message message) throws Exception {

    producer.sendOneway(message);

    }

    /**

  • 同一id消息放入同一个队列
  • */

    8.3.消费者

    1.创建consumer组

    2.设置Name Server地址

    3.设置消费位置,从最开始销毁

    4.设置消息回调处理监听 -> 处理消息

    public static void orderMessage(DefaultMQProducer producer,

    Message message, int id) throws Exception {

    // MessageQueueSelector:消息队列选择器,用于消息顺序发送

    SendResult sendResult = producer.send(message, new

    MessageQueueSelector() {

    @Override

    public MessageQueue select(List mqs, Message

    msg, Object arg) {

    //根据id(比如订单id)选择发送queue 这样就可以保证同一个id的消息

    会被发送到同一个queue中

    int index = id % mqs.size();

    return mqs.get(index);

    }

    }, id);//订单id

    System.out.printf("同消息放同队列发送成

    功,queueId:%d,topic:%s,tags:%s,内容:%s%n",

    sendResult.getMessageQueue().getQueueId(),

    message.getTopic(), message.getTags(), new String(message.getBody(),

    StandardCharsets.UTF_8));

    }

    /**

  • 延时消息
  • */

    public static void delayMessage(DefaultMQProducer producer,

    Message message) throws Exception {

    // 设置延时等级3,这个消息将在10s之后发送(现在只支持固定的几个时间,详看

    delayTimeLevel)

    // config/broker.conf中的配置

    // 默认配置 messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m

    7m 8m 9m 10m 20m 30m 1h 2h

    message.setDelayTimeLevel(3);

    SendResult sendResult = producer.send(message);

    System.out.printf("延时消息发送成

    功,queueId:%d,topic:%s,tags:%s%n",

    sendResult.getMessageQueue().getQueueId(),

    message.getTopic(), message.getTags());

    5.Start consumer

    // 消费者

    public class ConsumerTest {

    public static void main(String[] args) {

    try {

    /**

  • 实例化消息生产者Producer consumergroup
  • 1.每个消费组都会全量接收订阅的 topic和tag的消息,避免不同的消费者
  • 订阅topic和tag相同而使用的 consumerGroup不同导致的重复消费

  • 2.如果同一个消费组中的消费者订阅的 topic和tag是不相同的,后启动的
  • 消费者注册的topic和tag会覆盖原来注册到消费组,后面也只会接收对应的topic和tag的消息

  • 但消费组因为有多个消费者,所以就算消息的 topic和tag和消费者不
  • 匹配,仍然会轮询给 consumerGroup中的消费者,又因为该消费者订阅的 topic和tag不匹

    配,所以不会收到消息

  • 这种情况会造成消息的丢失
  • 3.如果topic相同tag不同,但消费逻辑一样,可以设置一个消费者订阅多
  • 个tag,实现消费者消费逻辑的通用

  • 消费逻辑不一样,就需要使用不同的consumerGroup
  • */

    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer

    ("consumergroup7");

    // 设置Consumer分批拉取消息的最大数 默认为1

    consumer.setConsumeMessageBatchMaxSize(1);

    //最大线程1个 默认最大线程 20,最小10个

    consumer.setConsumeThreadMax(1);

    consumer.setConsumeThreadMin(1);

    // 设置NameServer的地址

    consumer.setNamesrvAddr("127.0.0.1:9876");

    //从最开始的位置开始消费

    /**

  • CONSUME_FROM_LAST_OFFSET (默认值) 从队列的最后位置开始消费
  • 适用场景:新启动的消费者只想消费启动后新到达的消息
  • 注意:如果消费者组是第一次启动,且消息队列中没有任何消费进度记录,
  • 则从队列尾部开始消费(即跳过所有历史消息)

    *

  • CONSUME_FROM_FIRST_OFFSET 从队列的最前位置开始消费
  • 适用场景:需要处理队列中所有的历史消息(包括积压的)
  • 注意:这可能导致大量历史消息被重新消费,需评估系统负载能力
  • *

  • CONSUME_FROM_TIMESTAMP
  • 从指定时间点开始消费
  • 适用场景:需要消费某个特定时间点之后的消息
  • 8.4.消息的清理

  • 注意:需要配合 setConsumeTimestamp() 方法指定具体时间
  • *

  • 一旦消费者组已经有消费进度记录,此设置将不再生效,而是从服务端记录
  • 的进度继续消费

  • 重置消费进度可以使用 RocketMQ 提供的 resetOffsetByTime 等工具
  • */


    consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET

    );

    // 订阅一个或者多个Topic,以及Tag来过滤需要消费的消息和发送者保持一致

    才能搜到消息

    consumer.subscribe("topic_log7", "tags_error");

    // 注册回调实现类来处理从broker拉取回来的消息

    consumer.registerMessageListener(new

    MessageListenerConcurrently() {

    @Override

    public ConsumeConcurrentlyStatus

    consumeMessage(List msgs, ConsumeConcurrentlyContext

    context) {

    for (MessageExt msg : msgs) {

    String messageBody = new String(msg.getBody(),

    StandardCharsets.UTF_8);


    System.out.println(Thread.currentThread().getName() +"成功搜到消息,搜到消

    息数量:"+msgs.size()+","+"queueId:"+msg.getQueueId()+",搜到消

    息:"+messageBody);

    }

    //System.out.printf("%s 成功搜到消息: %s %n",

    Thread.currentThread().getName(), msgs);

    // 标记该消息已经被成功消费

    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;

    }

    });

    // 启动Producer实例

    consumer.start();

    System.out.printf("消费者1号启动成功.%n");

    }catch (Exception e){

    e.printStackTrace();

    }

    }

    }

    消息不会被单独清理,消息是顺序存储到commitlog的,消息是以commitlog为单位进行清理,R

    ocketMQ有

    自己的清理规则,默认是72小时候后进行清理

    ●到达时间清理点,自动清理过期的文件(凌晨4点)

    ●磁盘空间使用率达到了过期清理阈值(75%),自动清理过期的文件。

    ●磁盘占用率达到清理阈值(85%),开始按照设定的规则清理文件,从老的文件开始。

    ●磁盘占用率达到系统危险阈值(90%),拒绝写入数据。

    9.MQ消息的可靠性

    9.1 发送端消息可靠性:

    发送端Producer发送消息Broker端的核心逻辑如下图所示:

    RocketMQ架构模型中会有多个Borker为某个topic提供服务,一个topic下的消息分散存储在多个

    Broker存储端,它们是多对多关系。

    Broker会将其提供存储服务的topic的元数据信息上报到NameServer,对等NameServer节点组成

    的高可用服务会维护topic与Broker之间的映射关系,多对多的映射关系为消息可以重试发送到多

    个Broker端提供了前提与基础。

    当发送端需要发送消息时,如果发送端中缓存了topic的路由信息,并包含了消息队列,则直接返

    回该路由信息,如果没有缓存或没有消息队列,则向NameServer查询该topic的路由信息,查询到

    路由消息之后,采用指定的队列选择策略选择相应的queue发送消息,默认是采用轮询策略,发

    送成功则返回, 收到异常则根据相应的策略进行重试,可以根据发送端感知到的Broker的时延、上

    次发送失败的Broker信息和发送端配置的是否重试不同Broker的参数以及发送端设置的最大超时

    时间等等策略来灵活地实现不同等级的消息发送可靠性保证。

    重试策略可以有效的保证消息发送成功的概率,最终提高消息发送的可靠性。

    9.2 存储消息可靠性:

    RocketMQ的消息存储结构如下图所示:

    消息队列存储的最小单位是消息Message。

    同一个Topic下的消息映射成多个逻辑队列。

    不同Topic的消息按照到达broker的先后顺序以Append的方式添加至CommitLog,顺序写,随机

    读。

    目前RocketMQ存储模型使用本地磁盘进行存储,数据写入为producer -> direct memory -> pagec

    ache -> 磁盘,数据读取如果pagecache有数据则直接从pagecache读,否则需要先从磁盘加载到

    pagecache中。Broker存储节点的文件存储模式如下图所示:

    Broker端CommitLog采用顺序写,可以大大提高写入效率,同时采用不同的刷盘模式提供不同的

    数据可靠性保证,此外采用了ConsumeQueue中间结构来存储偏移量信息,实现消息的分发。由

    于ConsumeQueue结构固定且大小有限,在实际情况中,大部分的ConsumeQueue 能够被全部

    读入内存,可以达到内存读取的速度。此外为了保证CommitLog和ConsumeQueue的一致性, C

    ommitLog里存储了Consume Queues 、Message Key、Tag等所有信息,即使ConsumeQueue丢

    失,也可以通过 commitLog完全恢复出来,这样只要保证commitLog数据的可靠性,就可以保

    证Consume Queue的可靠性。RocketMQ存储端采用本地磁盘进行CommitLog消息数据的存储,

    不可避免的就会带来存储可靠性的挑战,如何保证消息不丢失,RocketMQ消息服务一直在不断

    提高数据的可靠性。

    9.3 Consume Queue的可靠性:

    RocketMQ存储端采用本地磁盘进行CommitLog消息数据的存储,不可避免的就会带来存储可靠

    性的挑战,

    如何保证消息不丢失,RocketMQ消息服务一直在不断提高数据的可靠性:

    1.存储可靠性挑战RocketMQ存储端也即Broker端在存储消息的时候会面临以下的存储可靠性挑

    战:

    ●Broker正常关闭

    ●Broker异常Crash

    ●OS Crash

    ●机器掉电,但是能立即恢复供电情况

    ●机器无法开机(可能是cpu、主板、内存等关键设备损坏)

    ●磁盘设备损坏

    1正常关闭:Broker 可以正常启动并恢复所有数据。

    2、3、4同步刷盘可以保证数据不丢失,异步刷盘可能导致少量数据丢失。

    5、6属于单点故障,且无法恢复。

    解决单点故障可以采用增加Slave节点,主从异步复制仍然可能有极少量数据丢失,同步复制可以

    完全避免单点问题。这里一般来说就需要在性能和可靠性之间做出取舍,对于RocketMQ来说,B

    roker的可靠性主要由两个方面保障:

    ● 单机的刷盘机制

    ● 主从之间的数据复制

    如果设置为每条消息都强制刷盘、主从复制,那么性能无疑会降低;如果不这样设置,就会有一

    定的可能性丢失消息。RocketMQ一般都是先把消息写到PageCache中,然后再持久化到磁盘

    上,数据从pagecache刷新到磁盘有两种方式,同步和异步。整体的消息写入和读取如下图所

    示:

    2. 同步刷盘:消息写入内存的 PageCache后,立刻通知刷盘线程刷盘,然后等待刷盘完成,刷

    盘线程执行完成后唤醒等待的线程,返回消息写成功的状态。这种方式可以保证数据绝对安

    全,但是吞吐量不大。

    3. 异步刷盘(默认):消息写入到内存的 PageCache中,就立刻给客户端返回写操作成功,当

    PageCache中的消息积累到一定的量时,触发一次写操作,或者定时等策略将 PageCache中

    的消息写入到磁盘中。这种方式吞吐量大,性能高,但是 PageCache中的数据可能丢失,不

    能保证数据绝对的安全。实际应用中要结合业务场景,合理设置刷盘方式,尤其是同步刷盘

    的方式,由于频繁的触发磁盘写动作,会明显降低性能。


    本文整理自实践经验,如有问题欢迎交流讨论。


    评论