Streams: a new general purpose data structure in Redis.

Salvatore Sanfilippo

Streams:Redis 中一种新的通用数据结构

直到几个月前,对我来说,streams 只是在消息传递语境下一个有趣且相对直观的概念。Kafka 让这个概念广为人知之后,我主要研究了它在 Disque 中的用途。Disque 是一个消息队列,现在计划将其转化为 Redis 4.2 module。后来我认为,Disque 主要关注的是 AP messaging,也就是容错能力和交付保证,而无需客户端付出太多努力,因此我决定 streams 这一概念并不适合这种场景。

然而与此同时,Redis 中存在一个问题,让我无法对默认提供的数据结构感到安心:Redis lists、sorted sets 和 Pub/Sub 的能力之间存在某种空白。你当然可以使用这些工具来建模消息或事件序列,但它们各自有不同的权衡。Sorted sets 很耗内存,无法自然地建模同一条消息被一次又一次传递的场景,客户端也无法阻塞等待新消息。因为 sorted set 不是一种顺序数据结构,而是一个元素可以通过改变 score 到处移动的集合:它不适合时间序列之类的场景,也就不足为奇了。Lists 在某些使用场景中也有类似的适用性问题:你无法探查 list 的中间部分,因为那里的访问时间是线性的。此外,list 无法实现 fan-out,针对 list 的阻塞操作只能将一个元素提供给一个客户端。Lists 中也没有固定的元素标识符,无法表达“给我从这个元素开始的内容”。对于一对多工作负载,可以使用 Pub/Sub,它在很多情况下都很出色;但在某些场景中,你并不想要 fire-and-forget:保留历史记录很重要,不仅是因为断开连接后需要重新获取消息,也因为某些消息列表(例如时间序列)非常适合通过范围查询来探索:这 10 秒范围内的温度读数是多少?

我解决上述问题的尝试,是计划将 sorted sets 和 lists 泛化为一种统一而更灵活的数据结构。然而,我的设计尝试几乎总是以这样的结果告终:所得数据结构比现有结构人工得多。Redis 的一件好事是,它对外提供的数据结构更像自然的计算机科学数据结构,而不是“Salvatore 发明的这个 API”。所以最后我停止了尝试,并说,好吧,目前我们只能提供这些;也许将来我会给 Pub/Sub 增加一些历史记录功能,或者让 lists 的访问模式更加灵活。然而,每当用户在会议期间走近我,问“在 Redis 中应该如何建模时间序列?”或类似的问题时,我的脸就变绿了。

Genesis(起源)

Redis 4.0 引入 modules 之后,用户开始思考如何自行解决这个问题。其中一位用户 Timothy Downs(蒂莫西·唐斯)通过 IRC 给我写了下面这段话:

<forkfork> 我计划开发的 module 是增加一种 transaction log style data type——也就是说,让大量订阅者能够实现类似 pub sub 的功能,同时不会让 redis 内存大幅增长

<forkfork> 订阅者在消息队列中保存自己的位置,而不是让 redis 维护每个消费者处理到了哪里,并为每个订阅者复制消息

这激发了我的想象。我思考了几天,意识到这可能正是我们一次性解决上述所有问题的时机。我需要重新构想“log”这一概念。它是一个基本的编程元素,人人都熟悉,因为它简单得就像以追加模式打开一个文件,然后按某种格式向其中写入数据。然而,Redis 数据结构必须是抽象的。它们位于内存中,而我们使用 RAM 不只是因为懒惰,更是因为借助少量指针,我们可以将数据结构概念化并使其抽象,从而摆脱那些显而易见的限制。例如,普通 log 通常有几个问题:offset 不是逻辑上的,而是实际的字节偏移;如果我们希望逻辑 offset 与条目插入的时间相关,该怎么办?这样我们就可以免费获得范围查询。同样,log 往往很难进行垃圾回收:如何从一种只追加的数据结构中删除旧元素?在我们理想化的 log 中,我们只需说,希望最多保留这个数量的条目,旧条目就会消失,诸如此类。

当我试图从 Timothy 的种子想法出发编写规范时,我正在开发一种 radix tree implementation,用于 Redis Cluster,以优化其内部的某些部分。这为实现一种非常节省空间、同时仍能以对数时间获取范围的 log 提供了基础。与此同时,我开始阅读 Kafka streams,以寻找适合我设计的其他有趣想法;这最终让我引入了 Kafka consumer groups 的概念,并针对 Redis 和内存使用场景再次对其进行了理想化。不过,这份规范连续几个月都只是一份规范;过了一段时间,我几乎从头重写了它,将与人们讨论这一即将加入 Redis 的功能时积累的许多建议融入其中。我希望 Redis streams 尤其适合时间序列这一用例,而不仅仅适用于其他类型的事件和消息传递应用。

让我们写点代码

从 Redis Conf 回来后,在夏季期间,我正在实现一个名为“listpack”的 library。这个 library 只是 ziplist.c 的后继者,也就是一种能够在单次分配中表示字符串元素列表的数据结构。它本质上是一种非常专门的序列化格式,特别之处在于它也可以按反向顺序解析,从右到左:这是为了在所有使用场景中替代 ziplists 所必需的能力。

将 radix trees 与 listpacks 混合使用,就可以很容易地构建出一种同时非常节省空间且带索引的 log,也就是说,允许通过 ID 和时间进行随机访问。准备好这些之后,我开始编写实现 stream 数据结构的代码。实现仍在收尾阶段,但此时 GitHub 上 Redis 的“streams”分支已经包含了足够多的内容,可以开始试用和玩耍了。我不敢说 API 已经 100% 定型,但有两个有趣的事实:第一,目前只缺 consumer groups,以及一些用于操作 stream、但重要性较低的命令;所有重要功能都已经实现。第二个事实是,一旦一切看起来稳定,我们决定大约两个月后将所有 stream 工作反向移植回 4.0 分支。这意味着 Redis 用户不必等待 Redis 4.2 才能使用 streams,它们会尽快提供给生产环境使用。这之所以可行,是因为作为一种新的数据结构,几乎所有代码改动都封装在新代码中。唯一的例外是 list 的阻塞操作:代码经过重构,使 streams 和 lists 的阻塞操作共享同一份代码,大幅简化了 Redis 的内部实现。

教程:欢迎使用 Redis Streams

在某种程度上,你可以把 streams 看作 Redis lists 的增强版。Stream 元素不只是单个字符串,而是由字段和值组成的对象。范围查询既可行又快速。Stream 中的每个条目都有一个 ID,它是一个逻辑 offset。不同客户端可以阻塞等待 ID 大于指定值的元素。Redis streams 的一个基本命令是 XADD。没错,所有 Redis stream 命令都以“X”作为前缀。

> XADD mystream * sensor-id 1234 temperature 10.5
1506871964177.0

XADD 命令会将指定条目追加为指定 stream“mystream”的新元素。上面的示例中,该条目有两个字段:sensor-id 和 temperature;不过同一 stream 中的每个条目都可以拥有不同字段。使用相同的字段名只会带来更好的内存利用率。还有一点很有意思:字段顺序保证会被保留。XADD 返回刚插入条目的 ID,因为第三个参数使用了星号,我们要求命令自动生成 ID。不过也可以强制指定一个 ID,例如为了将命令复制到 slaves 和 AOF 文件中。

ID 由两部分组成:毫秒级时间和序列号。1506871964177 是毫秒级时间,也就是精度为毫秒的 Unix 时间。点号后的数字 0 是序列号,用于区分在同一毫秒内添加的条目。这两个数字都是 64 位无符号整数。这意味着我们可以向 stream 中添加任意多条目,即使是在同一毫秒内也不例外。ID 的毫秒部分取 Redis server 生成 ID 时的当前本地时间与 stream 中最后一个条目之间的最大值。因此,即使计算机时钟向后跳跃,ID 仍会持续递增。在某种意义上,你可以把 stream 条目的 ID 看作完整的 128 位数字。不过,由于它们与添加条目的实例本地时间相关,我们无需额外付出代价就获得了毫秒精度的范围查询。

可以想见,以非常快的速度添加两个条目时,只有序列号会递增。我们可以简单地用一个 MULTI/EXEC block 来模拟“快速插入”:

> MULTI
OK
> XADD mystream * foo 10
QUEUED
> XADD mystream * bar 20
QUEUED
> EXEC
1) 1506872463535.0
2) 1506872463535.1

上面的示例还展示了如何为不同条目使用不同字段,而无需预先指定任何 schema。不过,实际情况是,每个 block 的第一条消息(通常包含约 50–150 条消息)都会被用作基准,后续字段相同的条目会通过一个标志进行压缩,该标志表示“与这个 block 的第一条条目使用相同字段”。因此,连续消息使用相同字段确实可以节省大量内存,即使字段集合会随时间缓慢变化也是如此。

从 stream 中获取数据有两种方式:由 XRANGE 命令实现的范围查询,以及由 XREAD 命令实现的 streaming。XRANGE 只会获取从 start 到 stop(包含两端)的条目范围。例如,如果我知道某个条目的 ID,就可以这样获取单个条目:

> XRANGE mystream 1506871964177.0 1506871964177.0
1) 1) 1506871964177.0
   2) 1) "sensor-id"
      2) "1234"
      3) "temperature"
      4) "10.5"

不过,你可以使用特殊的起始符号“-”和结束符号“+”,分别表示可能的最小和最大 ID。也可以使用 COUNT 选项限制返回的条目数量。下面是一个更复杂的 XRANGE 示例:

> XRANGE mystream - + COUNT 2
1) 1) 1506871964177.0
   2) 1) "sensor-id"
      2) "1234"
      3) "temperature"
      4) "10.5"
2) 1) 1506872463535.0
   2) 1) "foo"
      2) "10"

这里我们是按 ID 范围来处理的,但你也可以使用 XRANGE 获取给定时间范围内的特定条目范围,因为可以省略 ID 中的“sequence”部分。也就是说,你只需指定毫秒级时间。下面的命令表示:“从 Unix 时间 1506872463 开始,给我 10 条条目”:

127.0.0.1:6379> XRANGE mystream 1506872463000 + COUNT 10
1) 1) 1506872463535.0
   2) 1) "foo"
      2) "10"
2) 1) 1506872463535.1
   2) 1) "bar"
      2) "20"

关于 XRANGE,最后还有一点重要事项:由于回复中包含了 ID,而紧接着的下一个 ID 只需递增 ID 的序列部分即可轻松得到,因此可以使用 XRANGE 逐步遍历整个 stream,每次获取指定数量的元素。在 Redis 的 *SCAN 命令族之后,尽管 Redis 数据结构原本并不是为遍历而设计的,它们仍然支持遍历;我不想再犯同样的错误。

使用 XREAD 进行 streaming:阻塞等待新数据

当我们想按 ID 或时间获取 stream 范围,或按 ID 获取单个元素时,XRANGE 非常理想。但对于不同客户端必须随着数据到达而消费的 stream,这还不够好,需要某种形式的 pooling(对于某些只是偶尔连接来获取数据的应用,这可能是个好主意)。

XREAD 命令用于同时从多个 stream 中读取数据,只需指定我们在 stream 中已经获取到的最后一个条目的 ID 即可。此外,如果没有数据可用,我们可以要求阻塞客户端,并在数据到达时解除阻塞。这与 list 的阻塞操作类似,但这里数据不会从 stream 中被消费,多个客户端可以同时访问同一份数据。

下面是一个典型的 XREAD 调用:

> XREAD BLOCK 5000 STREAMS mystream otherstream $ $

它表示:从“mystream”和“otherstream”获取数据。如果没有数据可用,则将客户端阻塞 5000 毫秒,直到超时。在 STREAMS 选项之后,我们指定要监听的 key,以及已经获取到的最后一个 ID。不过,特殊 ID“$”表示:假设我已经拥有 stream 当前存在的所有元素,因此只需从下一个到达的元素开始提供。

如果我从另一个客户端发送以下命令:

> XADD otherstream * message "Hi There"

XREAD 端会发生如下情况:

1) 1) "otherstream"
   2) 1) 1) 1506935385635.0
         2) 1) "message"
            2) "Hi There"

我们会获得收到数据的 key,以及收到的数据。下一次调用中,我们很可能会使用最后收到的消息的 ID:

> XREAD BLOCK 5000 STREAMS mystream otherstream $ 1506935385635.0

以此类推。不过请注意,采用这种使用模式时,客户端可能会在很长时间之后才重新连接(因为处理消息需要时间,或出于任何其他原因)。在这种情况下,期间可能会积累大量消息,因此明智的做法是始终对 XREAD 使用 COUNT 选项,从而确保客户端不会被消息淹没,服务器也不必花费太多时间向单个客户端提供大量消息。

有上限的 streams

到目前为止一切顺利……不过,streams 迟早必须删除旧消息。幸运的是,XADD 命令提供的 MAXLEN 选项可以做到这一点:

> XADD mystream MAXLEN 1000000 * field1 value1 field2 value2

这基本上表示:如果添加新元素后发现 stream 包含超过 100 万条消息,就删除旧消息,使长度恢复到 100 万个元素。这就像对 lists 使用 RPUSH + LTRIM,只不过这次我们有了内置机制。不过请注意,上述操作意味着每次添加新消息时,我们还必须承担从 stream 另一端删除一条消息所需的工作。这会消耗一些 CPU,因此可以在 MAXLEN 的 count 前使用“~”符号,表示我们并不真正要求必须“恰好”有 100 万条消息,多几条也没什么问题:

> XADD mystream MAXLEN ~ 1000000 * foo bar

这样,XADD 只会在能够删除整个节点时才删除消息。与普通的 XADD 相比,这样就能几乎不花额外代价地维护有上限的 stream。

Consumer groups(开发中)

这是第一个尚未在 Redis 中实现、但正在开发中的功能。这一想法也最明显地受到了 Kafka 的启发,尽管在这里的实现方式相当不同。核心思路是,使用 XREAD 时,客户端还可以添加“GROUP <name>”选项。同一 group 中的所有客户端会自动获得不同的消息。当然,可以有多个 group 从同一个 stream 读取;在这种情况下,所有 group 都会收到 stream 中到达的同一消息的副本,但在每个 group 内,消息不会重复。

group 的一项扩展功能是:指定 group 时,可以加入“RETRY <milliseconds>”选项。在这种情况下,如果消息未通过 XACK 确认处理完成,那么在指定的毫秒数之后会再次传递这些消息。如果客户端没有私有手段标记消息已处理,这就能为消息传递提供某种尽力而为的可靠性。这部分功能也仍在开发中。

内存使用和保存、加载时间

由于 Redis streams 所采用的建模设计,其内存使用量非常低。具体取决于字段、值及其长度的数量,但对于简单消息,每使用 100 MB 内存就可以容纳数百万条消息。此外,这种格式的设计只需要极少的序列化:作为 radix tree 节点存储的 listpack block,在磁盘和内存中采用相同的表示形式,因此可以直接存储和读取。例如,Redis 可以在 0.3 秒内从 RDB 文件读取 500 万条条目。这使得 streams 的复制和持久化非常高效。

目前还计划允许删除中间的条目。该功能只实现了一部分,但策略是将条目标记为已删除;当条目与已删除条目之间达到给定比例时,重写 block 以回收垃圾,并在需要时将其与相邻的另一个 block 合并,以避免碎片化。

结论及预计时间

Redis streams 将在年底前成为 4.0 系列 Redis stable 的一部分。我认为,这种通用数据结构将为 Redis 打上一个巨大的补丁,使其能够覆盖大量过去难以覆盖的使用场景:也就是说,过去你必须发挥创造力,滥用现有数据结构来解决某些问题。一个非常重要的使用场景是时间序列,但我的感觉是,通过 TREAD 为其他使用场景进行消息 streaming 也会非常有趣:它既可以替代那些需要比 fire-and-forget 更高可靠性的 Pub/Sub 应用,也能支持全新的使用场景。目前,如果你想在自己的问题场景中开始评估这些新能力,只需从 GitHub 获取“streams”分支,然后开始试用。毕竟,欢迎提交 bug 报告 :-)

如果你喜欢视频,这里有一个实时展示 streams 的演示:https://www.youtube.com/watch?v=ELDzy9lCFHQ