Streams:Redis 全新的通用数据结构
原文由 Salvatore Sanfilippo 于 发布,订阅该博客
几个月前,对我而言,流(streams)不过是在消息传递语境下一个有趣且相对直观的概念。Kafka 让这一概念流行起来后,我主要是在 Disque 这一消息队列的场景下考察它的可用性,而 Disque 如今正计划被改写为 Redis 4.2 的一个模块。后来我认定,Disque 的核心是 AP 消息,也就是容错以及在客户端几乎无需额外付出的情况下保证消息送达,因此我认为流的概念并不适合那种场景。
不过,与此同时,Redis 上存在一个让我无法对默认提供的数据结构感到安心的问题。在 Redis 的 list、有序集合和 Pub/Sub 能力之间,存在某种空白。你当然可以用这些工具来建模一串消息或事件,但各有取舍。有序集合耗内存,无法自然地表达同一条消息被反复投递的情形,客户端也无法阻塞等待新消息。因为有序集合并非顺序型数据结构,它是一个集合,其中的元素可以通过修改分值被随意移动——难怪它不适合时间序列这类场景。List 则有另一类问题,在某些用例中同样会带来适用性上的限制:你无法探查列表中间的内容,因为那时的访问时间是线性的。而且无法实现扇出,针对 list 的阻塞操作一次只会把单个元素交给单个客户端。List 中也没有固定的元素标识,让你无法说“从这个元素开始给我后面的内容”。对于一对多的负载,有 Pub/Sub,它在很多情况下表现很好,但有些场景你并不想要“发后即忘”:保留历史很重要,不仅是为了在断线后重新获取消息,还因为某些消息序列——比如时间序列——非常需要通过范围查询来探查:比如“过去 10 秒内我的温度读数是多少?”
我曾尝试通过将有序集合和 list 泛化为一个更灵活的统一数据结构来解决上述问题,但我的设计尝试几乎总是让最终得到的数据结构比现有结构更加生硬、不自然。Redis 的一个优点在于,它对外提供的数据结构更像是自然的计算机科学数据结构,而不是“这是 Salvatore 发明的 API”。所以最后,我停下了尝试,心想,好吧,这就是我们目前能提供的,也许以后会给 Pub/Sub 加上历史记录,或让 list 的访问模式更灵活一些。然而,每当在会议上有人走过来问我“在 Redis 里该怎么建模时间序列?”或类似问题时,我的脸色就发青。
缘起
在 Redis 4.0 引入模块之后,用户开始尝试自行解决这个问题。其中一位用户 Timothy Downs 在 IRC 上给我写道:
<forkfork> 我打算做的模块是添加一种事务日志风格的数据类型——意思是大量订阅者可以实现类似 Pub/Sub 的功能,而不会导致 Redis 内存大幅增长
<forkfork> 由订阅者自己在消息队列中维护消费位置,而不是由 Redis 来跟踪每个消费者的进度并为每个订阅者复制消息
这段话点燃了我的想象。我思考了几天,意识到这或许正是能一次性解决上述所有问题的契机。我需要做的是重新构想“日志”这一概念。它是一个基础的编程元素,人人都很熟悉,因为它无非是以追加模式打开一个文件,然后以某种格式写入数据。然而 Redis 的数据结构必须是抽象的。它们驻留在内存中,我们使用内存并不只是图省事,而是因为借助若干指针,我们可以对数据结构进行概念化、加以抽象,从而摆脱显而易见的限制。举例来说,普通的日志有几个问题:偏移量不是逻辑上的,而是实际的字节偏移,如果我们想要与条目插入时间相关的逻辑偏移呢?这样就能免费获得范围查询。同样,日志往往难以进行垃圾回收:在一个只追加的数据结构中如何删除旧元素?而在我们理想化的日志中,只需说我们最多想要这么多条目,旧的就会自动消失,以此类推。
在我以 Timothy 的初步想法为起点撰写规范的同时,我正在为 Redis Cluster 开发一个基数树实现,用于优化其内部的某些部分。这为实现一个极其节省空间、同时仍能以对数时间进行范围访问的日志奠定了基础。与此同时,我开始阅读关于 Kafka streams 的资料,以寻找其他能很好融入设计的有趣想法,这让我引入了 Kafka 消费者组的概念,并再次针对 Redis 和内存场景对其进行了理想化重塑。不过,这份规范在数月里一直只是规范,以至于过了一段时间后,我几乎从头重写了一遍,以便融入与人们讨论这一即将加入 Redis 的新特性时积累的诸多启发。我希望 Redis streams 尤其能成为时间序列的绝佳适用场景,而不仅仅适用于其他类型的事件和消息应用。
开始写代码
从 Redis Conf 回来后的那个夏天,我正在实现一个名为“listpack”的库。这个库其实就是 ziplist.c 的继任者,也就是一种能在单次内存分配中表示字符串元素列表的数据结构。它只是一种非常专门的序列化格式,其特点是还能以逆序、从右到左的方式解析——这是要在所有用例中替代 ziplist 所必需的。
将基数树与 listpack 结合起来,就能轻松构建出一个既极其节省空间、又带有索引的日志,也就是说,可以通过 ID 和时间进行随机访问。这一基础就绪后,我便开始编写代码来实现 stream 数据结构。实现仍在收尾中,不过此刻在 GitHub 上的 Redis “streams” 分支里,已有足够的内容可以上手把玩了。我不敢说 API 已 100% 定型,但有两点值得关注:一是到目前为止,缺失的只有消费者组,以及若干用于操作 stream 的次要命令,所有核心功能都已实现;二是决定在大约两个月、一切看起来稳定之后,将全部 stream 相关工作反向移植回 4.0 分支。这意味着 Redis 用户不必等到 Redis 4.2 就能使用 streams,它们会尽快可用于生产环境。之所以能这样,是因为作为一种全新数据结构,几乎所有代码改动都自包含在新代码中。唯一的例外是阻塞式 list 操作:相关代码经过重构,使得 streams 和 list 的阻塞操作共享同一套代码,极大地简化了 Redis 内部实现。
教程:欢迎使用 Redis Streams
在某种程度上,你可以把 streams 看作 Redis list 的超级加强版。Stream 中的元素不再只是单个字符串,而是由字段和值组成的对象。范围查询是可行且高效的。Stream 中的每个条目都有一个 ID,即逻辑偏移量。不同的客户端可以阻塞等待 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,例如为了将命令复制到从节点和 AOF 文件。
ID 由两部分组成:毫秒级时间和序列号。1506871964177 是毫秒时间,其实就是毫秒精度的 Unix 时间。点号后面的数字 0 是序列号,用于区分在同一毫秒内添加的条目。这两个数字都是 64 位无符号整数。这意味着我们可以在一个 stream 中添加任意数量的条目,即使在同一毫秒内也是如此。ID 中的毫秒部分取值为生成该 ID 的 Redis 服务器当前本地时间与 stream 中最后一条条目时间二者中的较大值。因此,即便例如机器时钟回拨,ID 仍会保持递增。在某种意义上,你可以把 stream 条目的 ID 看作一个完整的 128 位数字。不过,由于它们与添加时所在实例的本地时间相关联,意味着我们免费获得了毫秒精度的范围查询能力。
正如你所料,以极快的速度添加两条条目,结果只会是序列号递增。我们只需通过一个 MULTI/EXEC 事务块就能模拟这种“快速插入”:
> MULTI OK > XADD mystream * foo 10 QUEUED > XADD mystream * bar 20 QUEUED > EXEC 1) 1506872463535.0 2) 1506872463535.1
上面的例子还表明,我们可以为不同条目使用不同的字段,而无需事先定义任何模式。不过实际发生的是,每个数据块(通常包含约 50 到 150 条消息)的第一条消息会被用作参照,后续具有相同字段的条目会通过一个标记进行压缩,该标记表示“与本块第一条条目的字段相同”。因此,为连续的消息使用相同字段确实能节省大量内存,即使字段集合会随时间缓慢变化也是如此。
要从 stream 中读取数据有两种方式:范围查询,由 XRANGE 命令实现;以及流式读取,由 XREAD 命令实现。XRANGE 只是按闭区间从起始到结束获取一段条目。例如,如果我知道某个 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 中的“序列号”部分。所以你只需指定毫秒级的时间即可。下面这条命令的意思是:“从 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 系列命令让原本并非为遍历而设计的数据结构也能被遍历之后,我避免了再次犯同样的错误。
通过 XREAD 进行流式读取:阻塞等待新数据
当我们想通过 ID 或时间获取范围、或通过 ID 获取单条元素来访问 stream 时,XRANGE 非常合适。但在多个客户端需要在数据到达时即刻消费 stream 的场景下,这还不够好,需要某种轮询(对于那些只是偶尔连接来获取数据的特定应用来说,轮询或许是个不错的主意)。
XREAD 命令被设计为可以同时从多个 stream 中读取,只需指定我们在该 stream 中已获取的最后一条条目的 ID。此外,如果没有可用数据,我们还可以请求阻塞等待,直到有数据到达再解除阻塞。这与阻塞式 list 操作的情形类似,但在这里数据并不会从 stream 中被消费,多个客户端可以同时访问同一份数据。
下面是 XREAD 调用的一个典型示例:
> XREAD BLOCK 5000 STREAMS mystream otherstream $ $
它的意思是:从“mystream”和“otherstream”获取数据。如果没有可用数据,则将客户端阻塞,超时时间为 5000 毫秒。在 STREAMS 选项之后,我们指定要监听的键以及我们已有的最后一条 ID。不过,特殊的 ID“$”表示:假定我已拥有该 stream 中当前存在的所有元素,因此只需从下一个到达的元素开始给我。
如果在另一个客户端中,我发送命令:
> XADD otherstream * message "Hi There"
这时在 XREAD 那一侧会发生如下情况:
1) 1) "otherstream"
2) 1) 1) 1506935385635.0
2) 1) "message"
2) "Hi There"我们会收到产生数据的键以及收到的数据。在下一次调用中,我们很可能会使用最后一条收到消息的 ID:
> XREAD BLOCK 5000 STREAMS mystream otherstream $ 1506935385635.0
依此类推。不过要注意,采用这种使用模式时,客户端可能会在很长一段时间后才再次连接(因为处理消息花了时间,或出于其他原因)。在这种情况下,期间可能会堆积大量消息,因此明智的做法是在使用 XREAD 时始终加上 COUNT 选项,以确保客户端不会被消息淹没,服务器也不必花费过多时间只为单个客户端提供大量消息。
带容量限制的 Stream
到目前为止一切顺利……不过 stream 迟早需要删除旧消息。好在通过 XADD 命令的 MAXLEN 选项就可以做到:
> XADD mystream MAXLEN 1000000 * field1 value1 field2 value2
这基本上意味着,如果在添加新元素后发现 stream 的消息数超过了 100 万,就删除旧消息,使长度回到 100 万条。这就像对 list 使用 RPUSH + LTRIM,只不过这次我们有了内置机制来实现。需要注意的是,上述做法意味着每次添加新消息时,还必须承担从 stream 另一端删除一条消息所需的工作。这会消耗一些 CPU,因此可以在 MAXLEN 的数量前使用“~”符号,以表示我们并非严格要求恰好 100 万条消息,多一点也无妨:
> XADD mystream MAXLEN ~ 1000000 * foo bar
这样,XADD 只会在能够删除整个节点时才去删除消息。相比普通的 XADD,这会让带容量限制的 stream 几乎不带来额外开销。
消费者组(正在开发中)
这是 Redis 中首个尚未实现、仍在开发中的特性。它也是受 Kafka 启发最为明显的想法,尽管在这里以相当不同的方式实现。要点在于,使用 XREAD 时,客户端还可以加上“GROUP <name>”选项。同一组内的所有客户端将自动获得不同的消息。当然,可以有多个组同时从同一个 stream 读取,在这种情况下,所有组都会收到 stream 中新到达消息的副本,但在每个组内部,消息不会重复。
对消费者组的一个扩展是,在指定组时还可以指定“RETRY <milliseconds>”选项:在这种情况下,如果消息未通过 XACK 确认已处理,就会在指定的毫秒数后再次投递。这在客户端没有私有方式来标记消息已处理时,为消息投递提供了某种尽力而为的可靠性。这一部分也仍在开发中。
内存占用与保存/加载时间
由于 Redis streams 所采用的设计,其内存占用非常低。具体取决于字段、值及其长度,但对于简单消息,每 100 MB 内存可容纳数百万条消息。此外,该格式被设计为只需极少的序列化:作为基数树节点存储的 listpack 块在磁盘和内存中具有相同的表示,因此存储和读取都非常简单。例如,Redis 能在 0.3 秒内从 RDB 文件中读取 500 万条条目。这使得 stream 的复制与持久化非常高效。
计划还将支持删除中间的条目。这部分目前只是部分实现,但策略是在条目标记中将条目标为已删除,当条目总数与已删除条目之间的比例达到一定阈值时,重写该块以回收垃圾,并在需要时将其与相邻的另一个块合并,以避免碎片化。
结论与预计时间
Redis streams 将在年底前作为 Redis 4.0 系列稳定版的一部分发布。我认为,这一通用数据结构将为 Redis 补上巨大的一块拼图,使其能够覆盖许多以往难以覆盖的用例:也就是说,过去你不得不绞尽脑汁、滥用现有数据结构来解决某些问题。一个非常重要的用例是时间序列,但我的感觉是,通过 TREAD 为其他用例进行消息的流式传输也会非常有吸引力,既可以作为对需要比“发后即忘”更高可靠性的 Pub/Sub 应用的替代,也能适用于全新的用例。目前,如果你想结合自身问题开始评估这些新能力,只需在 GitHub 上拉取“streams”分支并开始尝试。毕竟,我们欢迎提交 bug 报告 :-)
如果你喜欢看视频,这里有一个实时演示 streams 的视频:https://www.youtube.com/watch?v=ELDzy9lCFHQ
随机一篇博客
评论
登录后参与讨论