细说消息队列中
15、集群性能瓶颈和数据可靠性风险
消息队列的性能和可靠性由生产者、Broker集群、消费者三方共同保障,而不只是服务端的工作。通常衡量集群性能的一个重要指标是全链路耗时,即客户端发出一条消息到消费者消费到这条消息的时间差。
生产者的性能和可靠性
网络层面
网络层面,对性能和可靠性的影响主要包括连接协议、传输加密、网络稳定性、网络延时、网络带宽五个方面。
生产者客户端会先和Broker建立并保持TCP长连接,而不是在每次发送数据时都重新连接,以确保通信的性能。这也是默认情况下不用HTTP协议的原因。
在数据传输过程中,为了避免数据包被篡改、窃取,就需要进行传输加密。因为网络质量不稳定,传输过程中可能也存在丢包的情况,此时就需要依赖TCP的重传机制来解决问题。但当出现大量网络重传时,就会极大地影响性能,导致集群的吞吐下降和耗时上升。这也是我们在系统运营过程需要监控网络包重传率的原因。
消息发送模式
Broker的性能和可靠性
消费者的性能和可靠性
为了提高消息消费的及时性,最好是选择Push模型,即服务端有消息后主动Push给多个客户端,此时的消费的延时是最低的。从提高吞吐来看,为了避免服务端堆积,主流消息队列都是通过客户端主动批量Pull数据来提高吞吐、避免堆积。一般情况下,Pull模型都是默认的消费模型。
消息队列一般是通过消费分组(或订阅)消费数据,以便能自动分配消费关系和保存消费进度。此时当消费重平衡时,为了重新分配消费关系,所有的消费都会暂停,从而会影响到消费性能。如果重平衡次数较多,问题就会更加严重。所以,像Flink等流式计算引擎,都会绕过消费分组,指定分区进行消费,以避免重平衡带来的性能下降。而RocketMQ为了解决重平衡问题,就将重平衡移动到了Broker端,尽量降低消费重平衡带来的性能影响。
从可靠性来看,消费端是不存在丢数据的情况的。但是客户端如果存在错误提交消费位点(Offset)的情况,比如应该提交Offset却没有提交,就会导致重复消费;或者不应该提交Offset却提交了Offset,就会导致消费者没有消费到应该消费的数据,从而导致下游认为数据丢失。此时从代码上来看,建议是手动提交Offset(或ACK),即消费到数据,并且业务逻辑处理成功后,才执行ACK或者提交Offset。
16、分布式消息队列集群
有状态服务和无状态服务
我们经常听到有状态服务和无状态服务这两个词。二者之间最重要的一个区别在于:是否需要在本地存储持久化数据。需要在本地存储持久化数据的就是有状态服务,反之就是无状态服务。
HTTP Web服务就是典型的无状态服务。在搭建HTTP Web集群的时候,我们经常会使用Nginx或者在其他网关后面挂一批HTTP节点,此时后端的这批HTTP服务节点就是一套集群。
消息队列是有状态服务。消息是和分片绑定,分片是和节点绑定。所以,当需要发送一个消息后,就需要发送到固定的节点,如果把消息发送到错误的节点,就会失败。所以,为了将消息发送到对的节点和从对的节点削峰数据,消息队列在消息的收发上,就有服务端转发和客户端寻址两种方案。
消息队列的集群设计思路
元数据存储
消息队列集群元数据是指集群中Topic、分区、配置、节点、权限等信息。元数据必须保证可靠、高效地存储,不允许丢失。因为一旦元数据丢失,其实际的消息数据也会变得没有意义。
节点发现
节点探活
技术上一般分为主动上报和定时探测两种,这两种方式的主要区别在于心跳探活发起方的不同。从技术和实现上看,差别都不大,从稳定性来看,一般推荐主动上报。因为由中心组件主动发起探测,当节点较多时,中心组件可能会有性能瓶颈,所以目前业界主要的探活实现方式也是主动上报。
从探测策略上看,基本都是基于ping-pong的方式来完成探活。心跳发起方一般会根据一定的时间间隔发起心跳探测。如果保活组件一段时间没有接收到心跳或者主动心跳探测失败,就会剔除这个节点。比如每3秒探测一次,连续3次探测失败就剔除节点。探测行为一般会设置较短的超时时间,以便尽快完成探测。
主节点选举
节点启动
创建Topic
Leader切换
17、分布式的消息队列集群
元数据存储服务设计
基础第三方存储
一般只要具备可靠存储能力的组件都可以当作第三方引擎。简单的可以是单机维度的内存、文件,或者单机维度的数据库、KV存储,进一步可以是分布式的协调服务ZooKeeper、etcd等等。
集群内部自实现元数据存储
总结来看,在集群中实现元数据服务的优点是,后期架构会很简洁,不需要依赖第三方组件。缺点是需要自研实现,研发投入高。而如果使用独立的元数据服务,因为是现成的组件,产品成型就会很快,这也是当前主流消息队列都是依赖第三方组件来实现元数据存储的原因。所以当前主流消息队列的架构如下所示。
如图所示,Broker在启动或重连时,会根据配置中的ZooKeeper地址找到集群对应的ZooKeeper集群。然后会在ZooKeeper的/broker/ids目录中创建名称为自身BrokerID的临时节点,同时在节点中保存自身的Broker IP和ID等信息。当Broker宕机或异常时,TCP连接就会断开或者超时,此时临时节点就会被删除。
18、分布式集群可靠性
分区、副本和数据倾斜
副本间数据同步
ZooKeeper需要保证数据的高可靠,不允许丢失。而在多数原则理论中,如果数据只写入到Leader和Follower中,此时这两台节点同时损坏或者集群发生异常时导致Leader频繁切换,数据就可能会损坏或丢失。为了解决这些复杂场景,Zab协议定义了Zxid、崩溃恢复等细节来保证数据不会丢失。
19、Java分布式存储系统的编程技巧
PageCache调优和Direct IO
应用程序读取文件,会经过应用缓存、PageCache、DISK(硬盘)三层。即应用程序读取文件时,Linux内核会把从硬盘中读取的文件页面缓存在内存一段时间,这个文件缓存被称为PageCache。
可以绕过操作系统,直接使用通过自定义Cache + Direct IO来实现更细致、自定义的管理内存、命中和换页等操作,从而针对我们的业务场景来优化缓存策略,从而实现比PageCache更好的效果。从技术实现上看,它和我们使用C++管理内存编程是一样的效果。
FileChannel和mmap
FileChannel大多数时候是和ByteBuffer打交道的,你可以将ByteBuffer理解为一个byte[]的封装类。ByteBuffer是在应用内存中的,它和硬盘之间还隔着一层PageCache。
我们通过filechannel.write写入数据时,会将数据从应用内存写入到PageCache,此时便认为完成了落盘操作。但实际上,操作系统最终帮我们将PageCache的数据自动刷到了硬盘。这也是FileChannel提供了一个force()方法来通知操作系统进行及时刷盘的原因。
预分配文件、预初始化、池化
在提高文件写入性能的时候,预分配文件是一个简单实用的优化技巧。比如前面讲过,消息队列的数据文件都是需要分段的,所以在创建分段文件的时候,可以预先写入空数据(比如0)将文件预分配好。此时当我们真正写入业务数据的时候,速度就会快很多。
还一点就是对象池化,对象池化是指只要是需要反复new出来的东西都可以池化,以避免内存分配后再回收,造成额外的开销。Netty中的Recycler、RingBuffer中预先分配的对象都是按照这个池化的思路来实现的。
直接内存(堆外)和堆内内存
从应用程序的角度来看,可以通过批量同步刷盘的操作来提高性能。批量同步刷盘的核心思路是:每次刷盘尽量刷更多的数据到硬盘上。技术上是指通过收集多线程写过来的数据,汇总起来批量同步刷到硬盘中,从而提高数据同步刷盘的性能。
业务线程T1到T6的数据通过内存将数据发送给IO线程。然后业务线程进入await状态,当IO Thread收集到一定的数据后,再一起将数据同步刷到硬盘中。最后唤醒T1到T6线程,返回写入成功。
新的存储AEP
线程绑核
20、安全身份认证鉴权加密传输
消息队列的系统安全由六部分组成:网络隔离、传输安全、集群认证、资源授权、自我保护、数据加密。
网络隔离安全性
虚拟网络(VPC)就是一个独立的子网。如果不做对外打通(比如开通公网、跟其他网络拉专线),它的数据在这个网络内就是安全的。
数据传输加密
数据传输安全的核心是SSL/TLS,你可以简单理解成,如果要保证传输过程中的数据安全,就要用SSL/TLS。消息队列也是这个逻辑,几乎所有的消息队列产品,传输过程中的加密机制都是基于SSL/TLS实现的,你可以在它们的官方文档找到相应的资料,比如支持TLS的Pulsar或者RabbitMQ,支持SSL的Kafka。
SSL和TLS是同一个东西。SSL 3.0及之前的版本叫SSL,3.0之后叫做TLS,TLS是SSL的升级版。
建立连接时的身份认证
我们可以通过配置开启认证。认证是客户端连接服务端的第一道门槛,主要是解决客户端连接到服务端时,是否允许建立连接的问题。
认证就是客户端携带认证信息(比如用户名+密码、Token、Auth等)连接上服务端。服务端首先会判断当前开启或配置的认证类型,然后校验这些信息,如果通过,就允许建立连接,否则,就返回认证错误。
21、消息集群加密和限流方案
集群中的数据加密
当客户端发送消息A(比如hello world)到Broker,Broker保存到磁盘的数据是一串内容A加密后的字符串,当消费端消费数据时,Broker将加密后的字符串解密成消息A返回给消费端。
消息队列限流机制
全局限流方案
22、分布式监控设计
OpenTelemetry主要是解决可观测性数据的获取规范问题,类似消息队列领域的AMQP和OpenMessaging,目的都是打造一个标准化规范。它系统地将可观测性分为指标(Metrics)、日志(Logs)、跟踪(Traces)三个方面。在消息队列领域,可观测性建设主要也是围绕着这三点展开。
单机维度指标
消息队列关键指标
23、可观测性设计实现消息轨迹功能
消息唯一标识
一条消息的生命周期是从客户端发出来的时候开始算的,所以发送出来的时候,就应该赋予消息一个唯一ID,这才能真正表达消息的唯一。一旦在服务端赋予唯一ID,因为客户端可能会重复发送同一条数据,服务端就会认为是两条数据,生成两个ID。所以,消息ID的生成一般是在客户端SDK自动生成或者业务手动生成指定的。
消息轨迹设计原则
客户端轨迹数据记录
服务端轨迹数据记录
持久存储引擎的选择
24、RabbitMQ的集群架构
集群构建
数据可靠性
身份认证
资源鉴权
可观测性
25、RocketMQ的集群架构
集群构建
RocketMQ的元数据实际是存储在Broker上,不是直接存储在NameServer中。NameServer本身只是一个缓存服务,没有持久化存储的能力,先来看一张图示。
元数据信息实际存储在每台Broker上,每台Broker会在本节点维护持久化文件来存储元数据信息。这些元数据信息主要包括节点信息、节点上的Topic、分区信息等等。在Broker启动时,会先连接NameServer注册节点信息,并将保存的元数据上报到所有NameServer节点中。此时所有NameServer节点就有全量的元数据信息了,从而完成了节点之间的发现。
Broker和NameServer之间会有保活机制,Broker会定期和NameServer保持心跳探测,来确认节点运行正常。当Broker异常时,就会被踢出集群。
部署模式
因为是基于Raft算法实现的,所以根据Raft算法的多数原则,集群最少必须由三个节点来组成。不同节点的Raft Commitlog之间会根据Raft算法来完成数据同步和选主操作。当Master发生故障后,会先通过内部协商,然后从Slave节点中选出新的Master,从而完成主从切换。
因为实现方式的原因,Deledger模式最少需要三个节点,并且无法兼容RocketMQ原生的存储和复制能力(比如Master/Slave模式),而且这个模式维护较困难。所以为了解决这几个问题,RocketMQ将DLedger(Raft)能力进行上移,重新实现了选主组件DLedger Controller,这就是RocketMQ的Controller模式,也叫做DLedger Controller模式。
数据可靠性
Dledger模式副本间数据同步是采用同步写入的方式,即Master收到数据后,同步将数据写入到副本,多数副本写入成功后,就算数据写入成功。这个实现方式和ZooKeeper是一样的。从一致性上来看,这种方式属于最终一致。
安全控制
在传输安全方面,RocketMQ Broker支持TLS加密传输。从技术上看,RocketMQ Broker也是使用标准Java Server集成TLS的用法来实现的。
可观测性
接下来我们看看消息轨迹。RocketMQ的消息轨迹,在我看来是消息队列里面支持得最好的了。因为完整的消息轨迹需要包含生产者、Broker、消费者三部分的信息,如果需要支持生产端和消费端的轨迹信息,就需要在客户端SDK中集成轨迹信息上报的功能。RocketMQ在生产端和消费端实现了这个功能,而其他大部分消息队列在SDK是没有这个功能的。
如上图所示,RocketMQ的生产端和消费端的SDK集成了轨迹信息上报模块。当数据发送或消费成功时,如果开启轨迹上报,客户端会将轨迹数据上报到集群中的内置Topic或者自定义Topic中。因此Broker端就保存有全链路的轨迹信息了。
同时RocketMQ会为每条消息赋予一个唯一ID。当消息发送成功后,可以根据消息ID查看轨迹信息。如果需要,还可以把轨迹信息存储到一些第三方系统(比如Elasticsearch),以便后续查询。
26、Kafka的集群架构
数据可靠性
Kafka集群维度的数据可靠性也是通过副本来实现的,而副本间数据一致性是通过Kafka ISR协议来保证的。ISR协议是现有一致性协议的变种,它是参考业界主流的一致性协议,设计出来的符合流消息场景的一致性协议。
ISR协议的核心思想是:通过副本拉取Leader数据、动态维护可用副本集合、控制Leader切换和数据截断3个方面,来提高性能和可用性。
副本拉取Leader数据
Kakfa是通过Follower批量从Leader拉取数据来完成主从副本间的数据同步,并提高性能的。
动态维护可用副本集合
控制Leader切换和数据截断
在Kafka的实现中,这两种情况都是支持的,支持在Topic维度调整配置来选择这两个操作。所以Kafka的ISR协议的很大一部分工作,就是在代码层面处理Leader切换、数据阶段的操作。
安全控制
可观测性
27、Pulsar集群架构设计
集群构建
Pulsar Broker集群的构建思路和Kafka是一致的,都是通过ZooKeeper来完成节点发现和集群的元数据管理。
Broker启动时会在ZooKeeper上的对应目录创建名称为BrokerIP + Port的子节点,并在这个子节点上存储Broker相关信息,从而完成节点注册。
主节点选举
从技术上分析,是因为ZooKeeper在底层的存储数据结构是分层树结构。分层树结构在读取时需要多层检索,从而导致数据如果存储在硬盘,读取性能会很低。因此ZooKeeper只有将所有数据加载到内存中,才能提供较好的性能。
基于etcd的方案是当前集群化部署的推荐方案。etcd底层存储是B树的结构,在硬盘层面的读取性能较高,不一定要把数据加载到内存中,所以存储容量不受单机的限制。
基于RocksDB的方案和RabbitMQ的Mnesia大致上是一个思路,都是基于节点层面的存储引擎来完成元数据的存储。RocksDB的方案主要用在单机模式上,主要原因是RocksDB是一个单机数据库。
基于内存的方案主要用在集成测试的场景中。
数据可靠性
Pulsar是计算存储分离的架构,数据是通过Ledger和Entry的形式写入BookKeeper的。所以跟其他消息队列不一样的是,Pulsar的Topic没有副本概念,消息数据的可靠性是通过Ledger多副本来实现的。
Pulsar通过在Broker中设置Qw和Qa来设置Ledger的总副本数和写入成功的副本数。所以从一致性来看,Pulsar既可以是强一致,也可以是最终一致。
可观测性
需要获取完成文档,请扫下方二维码。
怎么获取大数据知识库
加入我们星球,即可获得超全大数据知识库和星球内容,内容超150万字,比市面上的大数据资料更全,比自己网上搜资料更靠谱,比你自学有更好的学习氛围,还有超多PDF文档供学习。
大数据知识库只是我们星球中的一部分内容,加入星球的权益还包括简历修改服务,精品大数据项目,问题答疑服务,大数据VIP专属微信群,精品资料分享,带领星友每日打卡学习... 加入我们星球,为你提供全方位服务。
在星球中也有超多的大数据内容(和上面大数据知识库中的内容不重复):
另外说明加入星球后支持三天无理由退款,不满意无条件随时退。 加入星球和语雀可享受以下服务:
一、 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
二、免费答疑解惑、交流技术
三、面试指导、模拟面试
四、各类pdf文档下载、星球代码下载
五、提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导
end
