Adventures in message queues

Salvatore Sanfilippo

消息队列探险记

原文由 Salvatore Sanfilippo 发布,订阅该博客

编辑注:如果你错过了,Disque 的源代码现已发布在 http://github.com/antirez/disque

过去几个月里,我花了大约 15% 到 20% 的时间,大多是从夜晚和周末挤出来的时间,在做一个新系统。它是一个消息代理,叫 Disque。我已经实现了最初设计里 80% 的功能,但仍然觉得还没到可以发布的程度。既然还发不了,那就至少写篇博客吧……这就是它如何起步的故事,以及关于它是什么的一些细节。

~ 初步探索 ~

很多开发者把 Redis 当作消息队列来用,往往是通过某个封装了 Redis 底层原语的库,有时则是直接用 Redis 原生 API 搭一个简单的临时队列。这种用法主要靠阻塞式的列表操作和列表写入操作来实现。说起来,用 Redis 干这个,似乎同时是最好的选择也是最糟的选择。好在它速度快、易于查看、部署和使用都很简单,而且在很多环境中它本来就是基础设施的一部分。但它也有缺点,因为 Redis 可变的数据结构与不可变的消息截然不同。Redis 在高可用/集群方面的权衡完全偏向于处理大的可变值,而同样的权衡用来处理消息却不是最佳选择。

对于消息代理来说,有一点至关重要,那就是要保证消息要么至少被投递一次,要么至多被投递一次。简单来说,鉴于要保证一条消息(这里的投递是指消息被 worker 接收*并*处理)恰好只被投递一次实际上是不可能的,选择就在于消息代理是能够保证 0 次或 1 次投递,还是 1 次到无限次投递。这通常被称为至多一次(at-most-once)语义和至少一次(at-least-once)语义。前者也有其用武之地,但更有意思、更实用的语义是后者,也就是保证消息至少被投递一次,如果出现故障就多次投递。

所以几个月前,我开始思考一种客户端协议,用一组 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 实现这一浩大工程,我采取了一种更务实的做法:我分叉了代码,并从源码中移除了所有 Redis 特有的部分,最终得到一个骨架。到这一步,我就可以开始实现我的设计了。

~ Disque 是什么? ~

经过几个月强度不大的工作和仅仅 200 次提交,我终于有了一个不再像玩具的系统:它有好几周看起来都像个玩具,以至于我都不敢谈论它,因为我随时可能把整个源码树删掉。现在大部分想法都已变成了带测试的可运行代码,我终于确信它将来会发布,也可以来谈谈我在设计中所做的权衡了。

Disque 默认就是一个分布式系统。由于它是一个 AP 系统,像 Redis 那样区分单节点模式和分布式模式就没有意义。单个 Disque 节点只是集群的一个特例,即只有一个节点的集群。所以这是设计中的一个要点:容错、抗分区,并且无论还有多少节点存活都保持可用,也就是所谓的 AP。我还希望这个系统能在不同场景下天然地扩展,无论是有大量生产者和消费者以及大量队列的场景,还是所有这些生产者和消费者都集中在单个队列上、而该队列可能分布在多个节点上的场景。

我的需求在大声告诉我一件事……Disque 将要做出一个重大的设计牺牲:消息顺序。Disque 只提供尽力而为的顺序保证。然而正因为这个牺牲,才有了很多收获……权衡很有意思,因为有时候它们会彻底打开设计的空间。

我可以就这样继续讲述 Disque 是什么,不过几个月前我在 Hacker News 上看到一条由 Jacques Chester 写的评论,见 https://news.ycombinator.com/item?id=8709146 [编辑:抱歉,我粘贴时弄错了 Adrian 的名字(嗨 Adrian,抱歉误引了你!)]。Jacques 恰好和我一样在 Pivotal 工作,他评论说不同的消息系统有着截然不同的功能和特性,如果没有细节,几乎不可能评估不同的选择,也无法判断一个系统比另一个快究竟是因为实现更好,还是仅仅因为它提供的保证少得多。于是他写了一组在评估消息系统时应该问的问题。我将借用他的问题,再补充几个,来描述 Disque 是什么,希望不是空泛而谈,而是能提供一些实际信息。

问:消息是否至少投递一次?

在 Disque 中,你可以选择至少一次投递(默认),或至多一次投递。这个属性可以按消息设置。至多一次投递只是至少一次投递的一个特例,即将消息的“retry”参数设为 0,并将消息只复制到单个节点。

问:消息是否由消费者确认?

是,消费者告知系统消息已正确投递的唯一方式就是确认它。

问:如果未被确认,消息会被多次投递吗?

是,Disque 会在“retry”时间之后自动重新投递消息,一直持续下去(直到消息达到最大 TTL)。当消息被确认后,确认信息会传播到持有该消息副本的节点。如果系统认为所有节点都已收到确认,消息最终会被垃圾回收并删除。已确认的消息在内存紧张时也会被驱逐。

节点会运行一种尽力而为的算法来避免多次入队同一条消息,以更好地逼近单次投递。然而在故障期间,多个节点可能会同时多次重新投递同一条消息。

问:队列是持久化的还是临时的?

持久化的。

问:持久化是通过先将每条消息写入磁盘,还是通过在服务器间复制消息来实现的?

默认情况下 Disque 仅在内存中运行,并通过同步复制来实现持久化(不过你可以按消息要求使用异步复制)。如果需要,也可以像 Redis 那样开启 AOF,比如在部署环境很可能出现大规模重启等情况时。在系统升级时,也可以仅为升级过程将 AOF 写入磁盘,以便即使平时不使用磁盘持久化,重启后也不会丢失状态。

问:队列是在一组服务器间部分/完全一致,还是为了最大化吞吐量而被分散的?

为了吞吐量而分散,不过消息顺序会以尽力而为的方式保留。每条消息都有一个不可变的“ctime”,它是一个以毫秒为单位的物理时钟时间戳,加上同一毫秒内生成消息的递增 ID。节点会利用这个 ctime 来对消息进行排序和投递。

问:在压力下消息会被完全丢弃吗?(即尽力而为模式)

不会,不过如果内存没有空间,新消息可能会被拒绝。当 75% 的内存已被使用时,接收消息的节点会尝试将消息外部复制到其他节点,而自己不保留副本,但如果其他节点也处于内存不足的状态,这种方式可能不会奏效。

问:消费者和生产者能否查看队列内部,还是队列完全不透明?

有可以“PEEK”查看队列的命令。

问:队列是无序的、先进先出(FIFO)的,还是带优先级的?

如前所述,是尽力而为的类 FIFO。

问:有代理(broker)还是无代理?

有代理,以一组主节点的形式存在。客户端可以与任意节点通信。

问:代理是否拥有独立的命名队列(主题、路由等),还是生产者和消费者需要协调彼此的连接?

命名队列。生产者和消费者无需协调,因为节点会通过联邦机制在集群内发现路由,并按消费者所需传递消息。不过,如果客户端愿意迁移到有更多消费者的地方,系统也会提供提示。

问:消息投递是事务性的吗?

是,一旦添加消息的命令返回,系统就会保证集群内已存在所需数量的副本。

问:消息接收是事务性的吗?

我想不算,因为如果未被确认,Disque 会尝试再次投递同一条消息。

问:消费者在接收时会阻塞,还是可以主动检查新消息?

两种行为都支持,默认是阻塞。

问:生产者在发送时会阻塞,还是可以检查队列是否已满?

如果本地节点中队列长度已经超过指定值,生产者可以在添加新消息时要求返回错误。

此外,如果生产者想尽快返回、让集群以尽力而为的方式复制消息,也可以要求异步复制消息。

没有办法在队列中消息过多时阻塞消费者,并在消息变少时再将其解除阻塞。

问:是否支持延时任务?

支持,以秒为粒度,最长可达数年。不过它们会占用内存。

问:消费者和生产者能否连接到不同的节点?

能。

希望通过这篇文章,Disque 能少一点“空中楼阁”的感觉。当然,不看代码很难判断,但如果最棒的特性已经公布,你至少已经可以开始吐槽了。以上有多少已经实现并能良好运行?除了 AOF 磁盘持久化,以及我想在 API 上再打磨的一些小细节之外,其余都已完成,所以首个版本应该不会太远,只是我做得很少,很难进展得非常快。

本文章由 muse-spark-1.2-contributor 进行翻译

评论