细说消息队列上
01、消息队列技术趋势
早年业界消息队列演进的主要推动力在于功能(如延迟消息、事务消息、顺序消息等)、场景(实时场景、大数据场景等)、分布式集群的支持等等。近几年,随着云原生架构和 Serverless 的普及,业界 MQ 主要向实时消息和流消息的融合架构、Serverless、Event、协议兼容等方面演进。从而实现计算、存储的弹性,实现集群的 Serverless 化。
02、业界主流消息队列
当时有个业务背景:用户的回复会包含图片和表情等复杂结构数据,导致内容可能会特别长。另外遇到高峰时,回帖数量会特别多,短时间内可能有几千万的回帖,回帖的总内容很大。
Redis 的问题在于数据都存在内存里,如果数据没有及时消费,就会打爆内存。而且因为回帖内容可能很大,在 Redis 的 Value 里存大的数据会有性能和稳定性问题。而 MySQL 在 insert 上是有性能瓶颈的,短时间内大量回帖,插入会特别频繁,性能就扛不住。
而 Kafka 就没有这个问题,作为一个消息队列,它的特点就是高吞吐、大消息、高并发、持久化,不会存在性能、功能、稳定性方面的问题。这就是我们选择 Kafka 的理由。
其实如果没有大消息和大流量等复杂场景,是可以选用非标准消息队列产品的。比如在用户状态审核的场景中,只需要向下游传递用户 ID 和审核结果,结构简单,数量有限。这时候选择非标准消息队列比如 Redis 和 MySQL 也是可以的。
所以,是否选择使用标准消息队列产品,取决于你的数据和业务场景的需求。当数据量大、场景复杂后,才必须引入标准消息队列,因为它有高吞吐、持久化、长久堆积的特性。
从宏观上来讲,我会认为具有缓冲作用、具备类发布和订阅能力的存储引擎都可以称做消息队列。因为消息队列的最基本功能就是生产和消费,在发布订阅之上,扩展如死信队列、顺序消息、延时消息等高阶能力,并实现高吞吐、低延时、高可靠等特性,就成为了我们所熟知的功能齐全的标准消息队列。
RabbitMQ 是 2007 年由国外一家叫做 Rabbit 的公司开源的,用 Erlang 写成的消息队列。主要满足业务中消息总线的场景。特点是功能丰富,低流量下稳定性较高,基本具备消息队列所应该具备的所有功能。缺点是在大流量的情况下会有明显的瓶颈和稳定性风险。
RocketMQ 在定位上和 RabbitMQ 很像,功能丰富,在业务消息中经常会用到。不过 RocketMQ 是在移动互联网浪潮下发展起来的,业务场景更加复杂,也支持更多功能,比如消息 Tag、消息轨迹、消息查询等等。
除了功能层面,在架构和性能层面,RabbitMQ 开发设计早,当时分布式的设计理念还不成熟,导致它在架构层面的设计存在较大的缺陷,遇到大流量、高并发的时候,容易出现集群不可用、网络分区等情况无法解决。而 RocketMQ 在分布式架构上实现得更合理优雅,在大流量、高并发的场景下表现优秀稳定。
Pulsar 和 Kafka 很像,主要定位在流领域,主打大吞吐的流式计算。但 Kafka 的功能比较简单,支持基本的发布订阅、幂等、事务消息。Pulsar 在满足这些功能的基础上,也希望支持 RocketMQ 和 RabbitMQ 的功能,所以功能最丰富。
除了功能层面,在架构和性能层面上,Pulsar 的架构设计比 Kafka 更符合当前云原生架构,它的定位是 Kafka 的升级版,主要解决 Kafka 当前的一些痛点问题,比如集群扩缩容慢、分区迁移需要 Rebalance、无法支持超多分区等。性能目前没有特别大的差距。


基于云原生架构设计的 Pulsar 开始走向成熟,业界的 MQ 也出现了计算存储分离、分层存储、多租户、弹性计算等概念。

03、消息队列的架构和功能
在系统架构中,消息队列的定位就是总线和管道,主要起到解耦上下游系统、数据缓存的作用。它不像数据库,会有很多计算、聚合、查询的逻辑,它的主要操作就是生产和消费。所以,我们在业务中不管是使用哪款消息队列,我们的核心操作永远是生产和消费数据。

下单流程是一个典型的系统解耦、消息分发的场景,一份数据需要被多个下游系统处理。另外一个经典场景就是日志采集流程,一般日志数据都很大,直接发到下游,下游系统可能会扛不住崩溃,所以会把数据先缓存到消息队列中。所以消息队列的基本特性就是高性能、高吞吐、低延时。
消息队列架构




消息队列功能




04、如何设计一个好的通信协议









设计的原则是:请求维度的通用信息放在协议头,消息维度的信息就放在协议体。那具体怎么设计呢?我们结合 Kafka 协议来分析。











gRPC 的生产消息的请求结构,内容简单,只需要定义生产请求需要携带的相关信息,比如 Topic、用户属性、系统属性等。
对比上面两段代码,我们可以很直观地看到区别,gRPC 的协议结构比 Remoting 的协议结构简单很多,比如不需要定义消息头,定义方式也更加简单。

05、消息队列的网络设计
网络模块的性能瓶颈分析



高性能网络模块的设计实现

因为语言特性、历史背景原因,RabbitMQ 用的就是这种方案。















而 Netty 就是这样一个基于 Java NIO 封装的成熟框架。所以我们一提到 Java 的网络编程,最先想到的就是 Netty。当前业界主流消息队列 RocketMQ、Pulsar 也都是基于 Netty 开发的网络模块,Kafka 因为历史原因是基于 Java NIO 实现的。
主流消息队列的网络模型实现






06、消息数据和元数据的存储
消息队列中的数据一般分为元数据和消息数据。元数据是指 Topic、Group、User、ACL、Config 等集群维度的资源数据信息,消息数据指客户端写入的用户的业务数据。
元数据信息的存储
元数据信息的特点是数据量比较小,不会经常读写,但是需要保证数据的强一致和高可靠,不允许出现数据的丢失。
同时,元数据信息一般需要通知到所有的 Broker 节点,Broker 会根据元数据信息执行具体的逻辑。比如创建 Topic 并生成元数据后,就需要通知对应的 Broker 执行创建分区、创建目录等操作。



另一种思路,集群内部实现元数据的存储是指在集群内部完成元数据的存储和分发。也就是在集群内部实现类似第三方组件一样的元数据服务,比如基于 Raft 协议实现内部的元数据存储模块或依赖一些内置的数据库。目前 Kafka 去 ZooKeeper 的版本、RabbitMQ 的 Mnesia、Kafka 的 C++ 版本 RedPanda 用的就是这个思路。


消息数据的存储


第一个思路,每个分区对应一个文件的形式去存储数据。具体实现时,每个分区上的数据顺序写到同一个磁盘文件中,数据的存储是连续的。因为消息队列在大部分情况下的读写是有序的,所以这种机制在读写性能上的表现是最高的。
但如果分区太多,会占用太多的系统 FD 资源,极端情况下有可能把节点的 FD 资源耗完,并且硬盘层面会出现大量的随机写情况,导致写入的性能下降很多,另外管理起来也相对复杂。Kafka 在存储数据的组织上用的就是这个思路。


第二种思路,每个节点上所有分区的数据都存储在同一个文件中,这种方案需要为每个分区维护一个对应的索引文件,索引文件里会记录每条消息在 File 里面的位置信息,以便快速定位到具体的消息内容。

因为所有文件都在一份文件上,管理简单,也不会占用过多的系统 FD 资源,单机上的数据写入都是顺序的,写入的性能会很高。缺点是同一个分区的数据一般会在文件中的不同位置,或者不同的文件段中,无法利用到顺序读的优势,读取的性能会受到影响,但是随着 SSD 技术的发展,随机读写的性能也越来越高。如果使用 SSD 或高性能 SSD,一定程度上可以缓解随机读写的性能损耗,但 SSD 的成本比机械硬盘高很多。




如果进行了分段,消息数据可能分布在不同的文件中。所以我们在读取数据的时候,就需要先定位消息数据在哪个文件中。为了满足这个需求,技术上一般有根据偏移量定位或根据索引定位两种思路。
根据偏移量(Offset)来定位消息在哪个分段文件中,是指通过记录每个数据段文件的起始偏移量、中止偏移量、消息的偏移量信息,来快速定位消息在哪个文件。
当消息数据存储时,通常会用一个自增的数值型数据(比如 Long)来表示这条数据在分区或 commitlog 中的位置,这个值就是消息的偏移量。



这两种方案所面临的场景不一样。根据偏移量定位数据,通常用在每个分区各自存储一份文件的场景;根据索引定位数据,通常用在所有分区的数据存储在同一份文件的场景。因为在前一种场景,每一份数据都属于同一个分区,那么通过位点来二分查找数据的效率是最高的。第二种场景,这一份数据属于多个不同分区,则通过二分查找来查找数据效率很低,用哈希查找效率是最高的。
消息数据存储格式

如果消息格式设计得不够精简,功能和性能都会大打折扣。比如冗余字段会增加分区的磁盘占用空间,使存储和网络开销变大,性能也会下降。如果缺少字段,则可能无法满足一些功能上的需要,导致无法实现某些功能,又或者是实现某些功能的成本较高。
所以,在数据的存储格式设计方面,内容的格式需要尽量完整且不要有太多冗余。


Kafka 的消息内容包含了业务会感知到的消息的 Header、Key、Value,还包含了时间戳、偏移量、协议版本、数据长度和大小、校验码等基础信息,最后还包含了压缩、事务、幂等 Kafka 业务相关的信息。


RocketMQ 的存储格式中也包含基础的 Properties(相当于 Kafka 中的 Header)、Value、时间戳、偏移量、协议版本、数据长度和大小、校验码等信息,还包含了系统标记、事务等 RocketMQ 特有的信息,另外还包含了数据来源和数据目标的节点信息。

消费完成执行 ACK 删除数据,技术上的实现思路一般是:当客户端成功消费数据后,回调服务端的 ACK 接口,告诉服务端数据已经消费成功,服务端就会标记删除该行数据,以确保消息不会被重复消费。ACK 的请求一般会有单条消息 ACK 和批量消息 ACK 两种形式。

因为消息队列的 ACK 一般是顺序的,如果前一条消息无法被正确处理并 ACK,就无法消费下一条数据,导致消费卡住。此时就需要死信队列的功能,把这条数据先写入到死信队列,等待后续的处理。然后 ACK 这条消息,确保消费正确进行。
这个方案,优点是不会出现重复消费,一条消息只会被消费一次。缺点是 ACK 成功后消息被删除,无法满足需要消息重放的场景。

这个方案,一条消息可以重复消费多次。不管有没有被成功消费,消息都会根据配置的时间规则或大小规则进行删除。优点是消息可以多次重放,适用于需要多次进行重放的场景。缺点是在某些情况下(比如客户端使用不当)会出现大量的重复消费。
ACK 机制和过期机制相结合的方案。实现核心逻辑跟方案二很像,但保留了 ACK 的概念,不过 ACK 是相对于 Group 概念的。
当消息完成后,在 Group 维度 ACK 消息,此时消息不会被删除,只是这个 Group 也不会再重复消费到这个消息,而新的 Group 可以重新消费订阅这些数据。所以在 Group 维度避免了重复消费的情况,也可以允许重复订阅。

纵观业界主流消息队列,三种方案都有在使用,RabbitMQ 选择的是第一个方案,Kafka 和 RocketMQ 选择的是第二种方案,Pulsar 选择的是第三种方案。

只有该段里面的数据都允许删除后,才会把数据删除。而删除该段数据中的某条数据时,会先对数据进行标记删除,比如在内存或 Backlog 文件中记录待删除数据,然后在消费的时候感知这个标记,这样就不会重复消费这些数据。
07、提升存储模块的性能和可靠性
内存读写的效率高于硬盘读写
批量读写的效率高于单条读写
顺序读写的效率高于随机读写
数据复制次数越多,效率越低
提升写入操作的性能
消息队列的数据最终是存储在文件中的,数据写入需要经过内存,最终才到硬盘,所以写入优化就得围绕内存和硬盘展开。写入性能的提高主要有缓存写、批量写、顺序写三个思路,这里对比来讲。

1. 缓存写和批量写
计算机多级存储模型的层级越高,代表速度越快(同时容量也越小,价格也越贵),也就是说写入速度从快到慢分别是:寄存器 > 缓存 > 主存 > 本地存储 > 远程存储。

内存读写的效率高于硬盘读写
批量读写的效率高于单条读写
写入优化的主要思路之一是:将数据写入到速度更快的内存中,等积攒了一批数据,再批量刷到硬盘中。


消息队列一般会同时提供:是否同步刷盘、刷盘的时间周期、刷盘的空间比例三个配置项,让业务根据需要调整自己的刷新策略。从性能的角度看,异步刷新肯定是性能最高的,同步刷新是可靠性最高的。
随机写和顺序写
单文件顺序写入硬盘很简单,硬盘控制器只需在连续的存储区域写入数据,对硬盘来讲,数据就是顺序写入的。

多文件顺序写入硬盘,系统中有很多文件同时写入,这个时候从硬盘的视角看,你会发现操作系统同时对多个不同的存储区域进行操作,硬盘控制器需要同时控制多个数据的写入,所以从硬盘的角度是随机写的。

提升写入操作的可靠性

一般消息队列都会开放这个配置项,默认批量刷盘,但有丢失数据的风险。如果业务需要修改为直接刷盘的策略来提高数据的可靠性,则会有一定的性能降低。

从理论来看,WAL 机制肯定会比直接写入缓存中的性能低。但我们实际落地的时候往往可以通过一些手段来优化,降低影响,达到性能要求。
虽然 WAL 日志需要极高的写入性能,但是数据量一般很小,而且是可顺序存储的、可预测的(根据配置的缓存大小和更新策略可明确计算)。
在实际落地中,我们可以采取 WAL 日志盘和实际数据盘分离的策略,提升 WAL 日志的写入速度。具体就是让 WAL 数据盘是高性能、低容量的数据盘,数据盘是性能较低、容量较大的数据盘,如果出现数据异常,就通过 WAL 日志进行数据
end
