消息队列历险记
编辑:如果你还没看到,Disque 的源代码现已在 http://github.com/antirez/disque 提供。
几个月来,我一直花费大约 15%~20% 的时间——主要是从夜晚和周末偷来的时间——开发一个新系统。它是一个 message broker(消息代理),名为 Disque。按照最初的规范,我已经实现了其中 80%,但我仍然觉得它还没准备好发布。既然还不能发布,那我至少可以写写博客……下面就讲讲它是如何开始的,以及它的一些细节。
~ 最初的几步 ~
许多开发者会把 Redis 当作消息队列使用,有时通过某个库进行封装,以抽象掉 Redis 的底层原语;有时则直接使用 Redis 的原始 API,构建简单的临时队列。这种用例主要通过阻塞式列表操作和列表推送操作来实现。Redis 显然同时是以这种方式使用时最好、也是最糟糕的系统。说它好,是因为它快速、易于检查、部署和使用,而且在许多环境中它本来就已经是基础设施的一部分。然而它也有缺点,因为 Redis 的可变数据结构与不可变消息非常不同。Redis 的 HA/Cluster 权衡完全偏向于大型可变值,但处理消息时,同样的权衡并不是最佳选择。
消息代理需要保证的一点是:一条消息要么至少投递一次,要么至多投递一次。简而言之,要保证消息恰好投递一次(这里的“投递”是指消息已被工作进程接收并处理)实际上是不可能的,因此可选方案是:消息代理能够保证投递 0 次或 1 次,或者投递 1 次到无限次。这通常被称为 at-most-once semantics(至多一次语义)和 at-least-once semantics(至少一次语义)。前者有其适用场景,但更有趣、更实用的是后者:保证消息至少投递一次,并在发生故障时多次投递。
因此,几个月前我开始思考一种客户端协议,尝试使用一组 Redis 主节点(完全不使用复制或集群),以提供这些保证。有时,只要针对某个用例稍微改变 Redis 的使用方式,就可能得到一个更好的系统。例如,在分布式锁方面,我曾尝试记录一种易于实现、但比单实例加故障转移实现更健壮的算法(http://redis.io/topics/distlock)。
然而,几天的工作之后,我的设计草案表明,设计一个临时系统可能是更好的选择,因为客户端算法最终过于复杂、并不理想,而我绝对想要实现的某些功能要么不可能实现,要么很难实现。向 Redis 添加更多功能似乎不是个好主意:它已经做了很多事情,而要做好消息传递,我需要的是与 Redis 运作方式截然不同的东西。但既然世界上到处都是消息代理,为什么还要设计一个新系统?因为数量惊人的用户正在使用 Redis,而不是专门为这一目标设计的系统,这很奇怪。少数人可能会错,但这么多人这样做,一定有某种原因。也许 Redis 的低门槛、简单 API 和速度,正是大多数人在接触消息代理领域时所习惯的东西。这个领域似乎充斥着两类方案:要么过于简单,把太多工作留给应用程序;要么过于复杂,但功能极其丰富。也许这里确实存在“消息领域的 Redis”这样的空间?
~ Redis 的彻底分叉 ~
人生中第一次,我没有一开始就直接写代码。几周里,我不时审视这个设计,把它转化为一个新系统,而不是 Redis 客户端库,并试着站在用户的角度理解:什么样的消息代理才能让我非常满意。最初的用例保持不变:延迟任务。Disque 是一个通用系统,但在设计过程中有 90% 的时间,我参考的都是这样一位用户:他需要解决发送消息的问题,而这些消息很可能是待处理的任务。如果某个东西不符合这个用例,我就把它删掉。
设计完成后,我终于开始编码。但该从哪里开始?“vi main.c”?幸好,Redis 在一定程度上是一个用 C 编写分布式系统的框架。我已经有了协议、网络库、客户端处理机制以及节点间消息总线。把这一切从头重写一遍,听起来实在是巨大的浪费。同时,我希望 Disque 在必要时能够在任何细节上与 Redis 完全分道扬镳,也希望它成为一个不会影响 Redis 本身的业余项目。因此,我没有尝试承担把 Redis 拆分成真正独立的框架和 Redis 实现这一庞大工作,而是采取了更务实的办法:我 fork 了代码,并从源代码中移除了所有 Redis 特有的内容,最终得到一个骨架。到此,我就可以实现自己的规范了。
~ Disque 是什么? ~
经过几个月非常断断续续的开发,以及仅仅 200 次提交,我终于有了一个不再像玩具的系统:它有好几个星期看起来都像个玩具,所以我甚至不敢谈论它,因为我直接删掉整个源代码树的可能性很大。现在,大部分想法都已经变成了带测试的可运行代码,我终于确信它将来会发布,也可以谈谈我在设计中做出的权衡。
Disque 默认是一个分布式系统。由于它是一个 AP system(AP 系统),所以没必要像 Redis 那样区分单节点模式和分布式模式。单个 Disque 节点只是集群的一种特殊情况,即集群中只有一个节点。因此,设计中的重要要点是:容错、能够抵御分区,并且无论还剩多少节点在线都能保持可用,也就是 AP。我还希望系统能够在不同场景下从根本上具备扩展能力:既能应对拥有许多队列的大量生产者和消费者,也能应对所有这些生产者和消费者都集中在单个队列上的情况,而这个队列可以分布到多个节点上。
我的需求大声告诉我一件事……Disque 将要在设计上做出一个巨大的牺牲:消息顺序。Disque 只提供 best-effort ordering(尽力而为的顺序保证)。然而正是由于这一牺牲,我们获得了很多东西……权衡之所以有趣,是因为它们有时会彻底拓展设计空间。
我本可以继续这样描述 Disque,但几个月前,我在 Hacker News 上看到了一条评论,作者是 Jacques Chester(雅克斯·切斯特),见 https://news.ycombinator.com/item?id=8709146[编辑:抱歉,我复制粘贴时弄错了 Adrian 的名字(Adrian,你好,很抱歉错误引用了你的话!)]。Jacques 恰好和我一样在 Pivotal 工作,他评论说,不同的消息系统拥有非常不同的功能集合和属性;如果不了解细节,几乎不可能评估不同的选择,也不可能判断一个系统比另一个更快,究竟是因为实现更好,还是仅仅因为它提供的保证少得多。因此,他写了一组问题,供人们在评估消息系统时提出。我会采用他的问题,再加上几个问题,用来描述 Disque 是什么,希望最终不是泛泛而谈,而是提供一些实际信息。
问:消息是否至少投递一次?
在 Disque 中,你可以选择 at-least-once delivery(至少一次投递,默认选项),也可以选择 at-most-once delivery(至多一次投递)。这个属性可以针对每条消息设置。将消息的“retry”参数设为 0,并将消息复制到单个节点,就可以实现 at-most-once delivery;它只是 at-least-once delivery 的一个特例。
问:消费者是否会确认消息?
会。消费者向系统表明消息已正确投递的唯一方式,就是确认消息。
问:如果消息没有得到确认,是否会多次投递?
会。经过一段“retry”时间后,Disque 会自动再次投递消息,并且会无限持续(但受消息最大 TTL 时间限制)。消息得到确认后,确认信息会传播到所有持有该消息副本的节点。如果系统认为所有节点都已收到确认,消息最终会被垃圾回收并删除。在内存压力下,已确认的消息也会被驱逐。
节点会运行一种 best-effort algorithm(尽力而为的算法),尽量避免将同一条消息排队多次,以便更好地近似单次投递。然而在发生故障时,多个节点可能会同时多次重新投递同一条消息。
问:排队是持久的还是临时的?
持久的。
问:持久性是通过先将每条消息写入磁盘,还是通过在服务器之间复制消息来实现的?
默认情况下,Disque 仅在内存中运行,并使用 synchronous replication(同步复制)实现持久性(不过你可以针对每条消息要求使用 asynchronous replication(异步复制))。如果预计部署环境可能发生大规模重启或类似情况,也可以按需启用 AOF(类似于 Redis)。系统升级时,可以仅为升级过程将 AOF 写入磁盘,从而即使平时不使用磁盘持久化,重启后也不会丢失状态。
问:在一组服务器中,排队是部分/完全一致的,还是为了最大吞吐量而进行分区?
为了吞吐量而进行分区,但消息顺序会以 best-effort 的方式得到保留。每条消息都有一个不可变的“ctime”,它由墙上时钟的毫秒级时间戳,以及同一毫秒内生成的消息所使用的递增 ID 组成。节点使用这个 ctime 对消息进行排序后再投递。
问:在压力下,消息可能被完全丢弃吗?(即 best effort)
不会。不过,如果内存没有空间,系统可能会拒绝新消息。当内存使用率达到 75% 时,接收消息的节点会尝试将消息复制到外部节点,即只复制到外围节点而不保留副本;但如果其他节点也处于内存不足状态,这种方式可能无法奏效。
问:消费者和生产者能否查看队列,还是队列完全不透明?
有用于“PEEK”队列的命令。
问:排队是无序的、FIFO,还是有优先级的?
如前所述,是尽力而为的、近似 FIFO 的顺序。
问:是否存在消息代理?
存在,由一组主节点组成。客户端可以与任意节点通信。
问:消息代理是否拥有独立的命名队列(主题、路由等),还是生产者和消费者需要协调各自的连接?
使用命名队列。生产者和消费者无需协调,因为节点会使用 federation(联邦机制)发现集群内部的路由,并根据消费者的需要传递消息。不过,如果客户端愿意将位置迁移到消费者更多的地方,系统会为它提供提示。
问:消息发布是否具有事务性?
是的。一旦添加消息的命令返回,系统就保证集群内部存在所需数量的副本。
问:消息接收是否具有事务性?
我猜不是,因为如果消息没有得到确认,Disque 会再次尝试投递同一条消息。
问:消费者在接收时会阻塞,还是可以检查是否有新消息?
两种行为都支持,默认情况下会阻塞。
问:生产者在发送时会阻塞,还是可以检查队列是否已满?
如果生产者要推送消息的本地节点上,消息长度已经超过指定值,生产者可以要求在添加新消息时返回错误。
此外,如果生产者希望尽快离开,并让集群以尽力而为的方式复制消息,也可以要求异步复制消息。
如果队列中的消息过多,没有办法阻塞消费者,并在消息减少后立即解除阻塞。
问:是否支持延迟任务?
支持,精度为秒,最长可达数年。不过它们会占用内存。
问:消费者和生产者能否连接到不同的节点?
可以。
希望通过这篇文章,Disque 不再那么像一纸空谈。当然,不查看代码很难判断,但既然你已经知道它最好的特性,至少可以开始抱怨了。上面说的内容有多少已经实现并且运行良好?除了 AOF 磁盘持久化,以及我想在 API 中进一步完善的一些小地方之外,其他都已经实现了。因此,首个版本应该不会太遥远,只是我投入开发的频率实在太低,很难进展得特别快。
随机一篇博客