Day13 RocketMQ完结

集群高级特性

  • DlDger文件一致性协议

    Dledger来自于开源组织OpenMessage

    • 为什么需要有一致性协议?

      在一个集群中,数据写入到server中的一个节点,然后希望从集群中任意一个节点都能读到写入的数据,实现这个功能有几个核心的问题

      1. 服务稳定性 各个server状态不稳定,随时可能宕机
      2. 网络抖动 server之间的网络如果发生抖动,请求就有可能丢失
      3. 网速问题 server之间的网络传输速度不一致,难以保证数据顺序
      4. 快速响应 在server之间相互同步数据的同时,还要快速给客户端响应结果 为解决这些问题,于是提出了一致性算法 弱一致性算法:DNS系统、Gossip协议 强一致性算法:Basic-Paxos、Multi-Paxos包括Raft(Nacos-JRaft、Kafka-KRaft、RocketMQ-Dledger)、ZAB等
    • 工作流程 image.png

      Log是保存在Server上的操作日志,其中的每个条目称为Entry,Entry中的操作,最终都会落到stateMachine中,Raft算法的核心就是要保证所有节点上的Entry顺序一致

      1. 多个server基于一致性协议,会共同选举产生一个Leader,负责响应客户端请求

      2. Leader通过一致性协议,将客户端的指令转发到集群所有节点上

      3. 每个节点将客户端的指令以Entry形式保存到自己的Log当中,此时Entry是uncommitted状态

      4. 当多数节点共同保存了Entry后,就可以执行Entry中的客户端操作,提交到State Machine中,此时Entry更新为commited状态

      Raft为每个节点设定了三种不同的角色 Leader:由选举产生,少数服从多数;向Follower节点发送心跳,Follower收到心跳就不会竞选Leader;响应客户端请求,集群内的所有数据变化都从Leader开始;向Follower同步操作日志 Follower:参与选举投票;同步来自Leader的数据;接收来自Leader/Follower的心跳,如果长时间没有接收到心跳,就自动转为Candidate Candidate:没有Leader时,发起投票竞选Leader 投票流程 所有节点起始角色都为Follower,每个节点都设定了一个随机选举过期时间Election Timeout(150ms-300ms),最早过期的那个节点会优先变成Candidate角色然后向其他节点发起投票请求,每个节点在同一任期Term内只有一次投票资格,在收到其他节点的投票响应后,Candidate会重置自身的Election Timeout,继续等待其他响应,一旦收获的票数超过集群节点数的一半,该节点就成为Leader节点,开始向其他节点发送心跳,确认自己的Leader身份,其他节点收到Leader心跳后,会自动转变为Follower状态,直到Leader心跳超时或者宕机才会触发下一次选举

  • 主从节点切换的高可用集群

    RocketMQ提供的Dledger虽然确实增加了集群的高可用,但是它是把集群选举和日志同步都一起完成的,而Dledger集群下的日志会比主从集群大很多,会增加写日志的IO负担,于是RocketMQ提供了一种Controller机制,既可以使用Raft选举机制,又可以使用原生的CommitLog日志

    image.png

    具体部署方式,参见官网 https://rocketmq.apache.org/zh/docs/deploymentOperations/03autofailover。

  • BrokerContainer容器式运行机制

    在RocketMQ4.x版本中,一个broker就是一个进程,而broker又是分主从的,压力方式不一样,主节点负责响应请求,而从节点一般就只承担冷备的作用,这种不对等的角色导致服务器资源不能充分利用。 于是在RocketMQ5.X版本中,提供了BrokerContainer模式。在一个BrokerContainer当中可以加入多个broker(master、slave、dledger),提高资源利用率,并且可以通过交叉部署来实现节点之间对等部署

    image.png

    bash
    复制代码
    # 修改配置文件 vim conf/container/broker-container.conf #配置端口,用于接收mqadmin命令 #listenPort=10811 #指定namesrv #namesrvAddr=worker1:9876;worker2:9876;worker3:98767 #或指定自动获取namesrv #fetchNamesrvAddrByAddressServer=false #指定要向BrokerContainer内添加的brokerConfig路径,多个config间用“:”分隔; #不指定则只启动BrokerConainer,具体broker可通过mqadmin工具添加 #brokerConfigPaths=/app/rocketmq/rocketmq-all-5.3.0-bin-release/conf/2m-2s-async/broker-b-s.properties:/app/rocketmq/rocketmq-all-5.3.0-bin-release/conf/2m-2s-async/broker-a.properties bin/mqbrokercontainer -c broker-container.conf

    image.png

RocketMQ5.x运行架构图

Kafka

Kafka服务搭建

  • 单机部署
bash
复制代码
# 1. 下载 Kafka(以 3.9.1 为例) cd /opt wget https://archive.apache.org/dist/kafka/3.9.1/kafka_2.13-3.9.1.tgz # 2. 解压并创建软链接 tar -zxvf kafka_2.13-3.9.1.tgz -C /usr/local/ ln -s /usr/local/kafka_2.13-3.9.1 /usr/local/kafka # 3. 配置环境变量 cat >> /etc/profile << 'EOF' export KAFKA_HOME=/usr/local/kafka export PATH=$PATH:$KAFKA_HOME/bin EOF source /etc/profile # 生成唯一集群 ID(保存输出值) kafka-storage.sh random-uuid # 示例输出:a9sqG57YRiW34orvtakVfg vim config/kraft/server.properties # 节点 ID(单机设为 1) # node.id=1 # 监听地址(替换为你的服务器 IP) # listeners=PLAINTEXT://:9092 # advertised.listeners=PLAINTEXT://<你的IP>:9092 # 日志存储目录(建议修改为非 /tmp 路径) # log.dirs=/var/log/kafka # 角色配置(单机同时作为 broker 和 controller) # process.roles=broker,controller # controller.listener.names=CONTROLLER # controller.quorum.voters=1@<你的IP>:9093 # 安全协议映射 # listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT # inter.broker.listener.name=PLAINTEXT # 使用上一步生成的集群 ID KAFKA_CLUSTER_ID="a9sqG57YRiW34orvtakVfg" kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /usr/local/kafka/config/kraft/server.properties # 输出提示:Formatting /data/kafka-logs with metadata.version 3.9-IV0 # 后台启动 Kafka nohup kafka-server-start.sh /usr/local/kafka/config/kraft/server.properties > /var/log/kafka/kafka.log 2>&1 & # 验证进程 ps -ef | grep kafka # 应看到:kafka.Kafka config/kraft/server.properties # 1. 创建测试主题 kafka-topics.sh --create --topic test-topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1 # 2. 查看主题列表 kafka-topics.sh --list --bootstrap-server localhost:9092 # 3. 发送消息(在新终端) kafka-console-producer.sh --bootstrap-server 60.204.184.61:9092 --topic test-topic # 输入消息内容后按 Enter # 4. 消费消息(在新终端) kafka-console-consumer.sh --bootstrap-server 60.204.184.61:9092 --topic test-topic --from-beginning

常见MQ问题梳理

  • 如何保证消息不丢失?

    首先要明确在使用MQ的链路中,哪些场景可能丢失消息?

    1. 生产者发到Server时,有可能因为网络原因造成数据丢失
    2. Server主从同步时,有可能因为网络造成数据丢失
    3. 消息数据刷盘时(可能性较小),数据由page cache写入磁盘时,服务宕机,有可能造成数据丢失
    4. Server推送到消费者时,有可能因为网络造成数据丢失 image.png
    1. 生产者发送消息如何保证不丢失?

      生产者发送消息确认机制

      1. RocketMQ实现有三种发消息的方式 producer.sendOneWay(msg) 异步发送,不确认是否成功,有丢失消息的风险 producer.send(msg) 同步发送,等待broker确认已收到消息再进行下一步处理 producer.send(msg, new SendCallback()) 异步发送,回调确认消息是否发送成功

      也可以使用RocketMQ的事物消息,确保与数据库事物修改同步 2. Kafka实现Future future = producer.send(record) ,生产者调用future.get()方法拿到元数据,就说明已成功 3. RabbitMQ提供了Publisher Confirms机制,

      //添加两个回调,一个处理ack响应,一个处理nack响应

      channel.addConfirmListener(ConfirmCallback ackCallback, ConfirmCallback nackCallback);

    2. Broker如何保证消息不丢失?操作系统刷盘操作,没法完全保证数据安全性

      1. RocketMQ 提供了一个刷盘配置flushDiskType,SYNC_FLUSH同步刷盘,发一次消息就写一次磁盘(其实底层代码是间隔10ms),ASYNC_FLUSH异步刷盘,间隔一段固定时间刷盘,性能更稳定

      2. Kafka提供了参数可设置刷盘策略,当log.flush.interval.messages为1时,就是一条消息刷一次盘

        bash
        复制代码
        flush.ms : 多长时间进行一次强制刷盘。 log.flush.interval.messages:表示当同一个Partiton的消息条数积累到这个数量时,就会申请一次刷盘操作。默认是Long.MAX。 log.flush.interval.ms:当一个消息在内存中保留的时间,达到这个数量时,就会申请一次刷盘操作。他的默认值是空。如果这个参数配置为空,则生效的是下一个参数。 log.flush.scheduler.interval.ms:检查是否有日志文件需要进行刷盘的频率。默认也是Long.MAX。
      3. RabbitMQ明确说明即便是持久化队列,也不会每条数据都刷盘,基本是交由操作系统自行处理

    3. Broker主从同步怎么保证消息不丢失?

      1. RocketMQ

        • 普通集群 由于RocketMQ普通主从集群主节点挂了之后不会自动切换主节点,未同步的数据仍旧会保留在主节点上,待主节点恢复重启后,从节点会将未同步到的消息继续同步,但也正因为如此,主节点崩溃(Q5)将没有新的消息可写入
        • Dledger集群 采用CP模型,将数据在主从节点之间同步,两阶段提交CommitLog确保多数节点能够同步到消息,即使切换主节点,也能极大程度的保留数据,但是极端情况下仍旧会造成数据丢失,不过可能性非常小
      2. Kafka

        在Kafka集群采用的是AP模型,在集群中,如果Leader Partition服务宕机,Follower会重新选举产生一个新的Leader Partition,而所有的消息都会以Leader为准,这样即使旧的Leader重启了,也是作为Follower将保存的那份数据刷掉

    4. 消费者消费怎么保证消息不丢失?

      消费状态确认机制 消费者处理完成后,需要给Broker一个响应,告诉对方自己已经把消息妥善处理了,如果broker没有收到响应,就会认为消息没有处理成功,会根据重试机制再次尝试消息投递,RocketMQ和Kafka根据offset进行重新投递,而RabbitMQ的Classic Queue会将消息重新入队,正常来说,这个过程丢消息的可能性极小,消费者应该考虑重复投递的幂等性问题

    5. 如果MQ服务挂了,怎么保证消息不丢失?

      设计缓存降级,在进行消息投递重试失败后,考虑将消息往降级缓存队列里写,继续正常执行业务,然后不断将降级缓存里的数据投递到MQ服务,等MQ服务恢复后流程就正常了

  • 怎么保证消息的顺序性?

    由于队列的FIFO性质,在同一个Message Queue中,消息是天然有序的,那么要保证整体消息的有序性(一般是局部有序),就应该保障两个方面

    1. Producer 将一组有序数据写入到同一个Message Queue

      RocketMQ和Kafka都提供了分区机制,可以让应用程序自行决定写到哪个分区(Message Queue),即能保证同一批数据的写入有序性;

      而RabbitMQ则可以通过维护Exchange与Queue之间的绑定关系,将这一组局部有序的消息转发到同一个对列中,从而保证这一组有序的消息,

    2. Consumer每次集中从同一个Message Queue中拿取消息

      RocketMQ 消费是多线程的,但可以提供实现MessageListenerOrderly接口来对消费线程进行并发控制,保证消费有序性 Kafka 消费某一个partition时,天生就是单线程,所以也能过保证有序性 RabbitMQ 本身就是点对点的消费模式,只要保证一个队列中只有一个consumer,就能是顺序的,如果某条消息消费失败,由于它会将消费失败的数据进行重新投递,那就无法保证顺序性了

  • 如何保证消息幂等性?

    1. 生产者投递过程中的幂等性

      RocketMQ 在发送消息时给每条消息分配一个唯一ID,提供这个ID判断是否重复投递

      Kafka 提供幂等性配置,需要打开idempotence幂等性控制(默认是打开的,但是如果其他配置有冲突,会影响幂等性配置) 生产者投递消息不会使用MQ系统自带的处理,而是在消息本身带有唯一ID或者版本号,然后由消费者自行做幂等判断

    2. 消费者消费过程中的幂等性

      消费时,校验唯一ID或版本号判断是否继续消费,除此之外,需要关注重试队列和死信队列中的数据,及时排查消费失败的原因

  • 如何快速处理消息积压?

    消息积压的危害 当MQ有大量消息积压时,会严重影响服务端的性能。 对于RocketMQ和Kafka来说,消息积压能力很强,短时间积压不会有太大影响,但是如果一直得不到解决,积压数据的日志文件可能丢失。 对于RabbitMQ来说,StreamQueue流式队列处理机制类似于以上两者,Classic和Quorum处理能力弱,需及时处理积压问题

    消息积压的根本在于消费者consumer处理效率太低,所以最核心的目标就是提升consumer消费效率,消费效率上来了问题自然迎刃而解。

    对于RabbitMQ,由于是点对点模式,一条消息只能由一个consumer消费,所以只需增加consumer的数量,即可解决。

    对于RocketMQ和Kafka,由于同一个消费者组下的多个Consumer需要和对应topic下的MessageQueue建立对应关系,而一个MessageQueue最多只能被一个消费者consumer消费,因此,即使增加消费者,至多也只能和topic下MessageQueue的数量相同,如果此时MessageQueue的数量本身不够多的话,也就没法使用这种方法解决了。此时可以创建一个新的临时topic,配置足量的MessageQueue,然后把Consumer实例的topic指向新topic,并紧急上线一组新的topic去消费旧topic的数据往新topic里面丢(只处理这一个逻辑),这样处理能够快很多,但也是紧急使用,后续要对原有的设计进行优化。

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP