跳至主内容

流过滤内部机制

·10 分钟阅读

之前的帖子介绍了 RabbitMQ 3.13 中的一项令人兴奋的新功能——流过滤。在本帖中,我们将介绍流过滤的内部机制。了解其设计和实现将有助于您以最适合您的用例的方式来配置和使用流过滤。

概念

流过滤的核心理念是在代理(Broker)端提供第一层高效过滤,且无需代理去解析消息内容。这样,只需要流中部分数据的消费者就不必获取所有数据并自行处理过滤。这可以大幅减少传输给消费者的数据量。

通过过滤,每个消息可以关联一个过滤值(filter value)。它可以是地理信息,例如消息来源的区域,如下图所示:

Each message can have a filter value associated to it, like the region of the world it comes from.
每条消息都可以关联一个过滤值,例如其来源区域。

因此,我们的流中有 1 条 AMER(绿色)消息,1 条 APAC(深蓝色)消息,2 条 EMEA(紫色)消息,然后是 2 条 AMER 消息。

消息发布

发布者负责为每条外发消息关联其过滤值。

A publisher provides the filter value for each outbound message.
发布者为每条外发消息提供过滤值。

在上图中,发布者发布了 1 条 AMER(绿色)消息和 2 条 EMEA(紫色)消息,这些消息将被添加到流中。

消息消费

当消费者进行订阅时,它可以指定一个或多个过滤值,代理将只分发带有这些过滤值的消息。我们稍后会看到实际情况略有不同,但理解这一概念就足够了。

在下图中,顶部的消费者指定它只需要 AMER(绿色)消息,代理就只分发这些消息。中间的 EMEA 消息消费者和底部的 APAC 消息消费者也是如此。

A consumer can specify it only wants messages with given filter value(s) and the broker will do its best to dispatch messages accordingly.
消费者可以指定它只需要带有特定过滤值的消息,代理将尽最大努力进行相应的分发。

以上就是相关概念,现在让我们来探索一下实现细节。

流的结构

我们需要了解流的结构才能理解流过滤的内部机制。流是一个包含分段文件(segment files)的目录。每个分段文件都有一个关联的索引文件(用于确定在分段文件的给定偏移量处挂载消费者等)。拥有多个“小型”分段文件比拥有一个巨大的单体文件要好得多:例如,删除“旧”分段文件来截断流比移除大文件的开头部分更高效、更安全。

分段文件由包含消息的“数据块”(chunks)组成。一个数据块中的消息数量取决于进入速率(高进入速率意味着一个数据块中包含许多消息,低进入速率则包含较少消息)。一个数据块中的消息数量从几条(甚至 1 条)到几千条不等。

数据块(Chunks)有什么作用?数据块是流中的工作单元:它们用于复制,更重要的是对于我们的主题而言,它们用于消费者投递。代理使用 sendfile 系统调用将数据块发送给消费者(将整个数据块从文件系统直接发送到网络套接字,而无需将数据拷贝到用户空间)。

下图说明了流的结构:

A stream is made of segment files. A segment file contains chunks, a chunk contains messages. The broker dispatches chunks to consumers, not individual messages.
流由分段文件组成。分段文件包含数据块,数据块包含消息。代理分发的是数据块,而不是单独的消息。

有了这些基础,让我们看看代理如何判断是否要分发某个数据块。

代理端过滤

假设有一个只想要 AMER(绿色)消息的消费者。当代理准备分发一个数据块时,它需要知道该数据块是否包含 AMER 消息。如果包含,它可以发送该数据块;如果不包含,代理可以跳过该数据块,移动到下一个并重复此过程。

每个数据块都有一个头部,其中可以包含一个布隆过滤器(Bloom filter),它告诉代理该数据块是否包含带有特定过滤值的消息。布隆过滤器一种空间效率极高的概率型数据结构,用于测试一个元素是否属于某个集合。在我们的例子中,集合包含 AMEREMEAAPAC,而目标元素是 AMER

下图说明了针对 3 个数据块的代理端过滤过程:

Filtering on the broker side. The broker uses a Bloom filter contained in each chunk to know whether it contains messages with a given filter value. A Bloom filter is efficient but it can return false positives, so the broker may return chunks that do not contain messages with the requested filter value(s).
代理端过滤。代理使用每个数据块中包含的布隆过滤器来判断其是否包含特定过滤值的消息。布隆过滤器很高效,但可能会产生误报,因此代理可能会返回不包含请求过滤值的消息的数据块。

如上图所示,过滤器可能会产生误报,即返回不包含预期过滤值消息的数据块。这是正常的,因为布隆过滤器是概率性的。不过它们不会产生漏报:如果过滤器说没有 AMER(绿色)消息,那我们就可以确定是真的没有。我们必须接受这种不确定性:有时我们可能会白白分发一些数据块,但这总比分发所有数据块要好。

可以肯定的是,消费者确实可能会收到它不想要的消息:看左边的第一个数据块,它包含了消费者要求的 AMER(绿色)消息,但也包含了 EMEA(紫色)和 APAC(深蓝色)消息。这就是为什么客户端也必须进行过滤的原因。

客户端过滤

代理在投递消息时处理了第一层过滤,但由于投递单元是数据块,消费者仍然可能收到不需要的消息。因此,客户端也必须进行过滤,这显然必须与订阅时设置的过滤值保持一致。

下图说明了一个只想要 AMER(绿色)消息的消费者,它必须执行最后一步过滤:

As a chunk can contain unwanted messages, the client must filter out messages as well.
由于数据块可能包含不想要的消息,客户端也必须过滤掉这些消息。

让我们看看这如何转化为应用程序代码。

API 示例

过滤不是侵入性的,可以作为横切关注点来处理,从而最小化对应用程序代码的影响。以下是在使用 Java 流客户端声明生产者时,如何设置从消息中提取过滤值的逻辑(filterValue(Function<Message,String>) 方法):

Producer producer = environment.producerBuilder()
.stream("invoices")
.filterValue(msg -> msg.getApplicationProperties().get("region").toString())
.build();

在消费端,Java 流客户端提供了 filter().values(String... filterValues) 方法来设置过滤值,并提供了 filter().postFilter(Predicate<Message> filter) 方法来设置客户端过滤逻辑。在声明消费者时必须调用这两个方法。

Consumer consumer = environment.consumerBuilder()
.stream("invoices")
.filter()
.values("AMER")
.postFilter(msg -> "AMER".equals(msg.getApplicationProperties().get("region")))
.builder()
.messageHandler((ctx, msg) -> {
// message processing code
})
.build();

如你所见,过滤不会改变发布和消费代码,只会影响生产者和消费者的声明。

现在让我们看看如何针对特定用例以最合适的方式配置流过滤。

流过滤配置

关于流过滤的第一篇文章提供了一些数据(与无过滤相比,使用过滤可节省约 80% 的带宽)。流过滤的收益很大程度上取决于用例:进入速率、过滤值的基数和分布,以及过滤器大小。过滤器越大越好(错误率越低)。在数据块中使用的过滤器大小可以设置为 16 到 255 字节之间,默认值为 16 字节。

Java 流客户端提供了 filterSize(int) 方法,用于在创建流时设置过滤器大小(它在内部设置 stream-filter-size-bytes 参数)。

environment.streamCreator()
.stream("invoices")
.filterSize(32)
.create()

如何估算过滤器的大小?网上有很多布隆过滤器计算器。其参数包括哈希函数数量(RabbitMQ 流过滤为 2)、预期元素数量、错误率和大小。通常你对元素数量有个大致概念,因此需要在错误率和过滤器大小之间找到平衡点。

以下是一些示例:

  • 10 个值,16 字节 => 2% 错误率
  • 30 个值,16 字节 => 14% 错误率
  • 200 个值,128 字节 => 10% 错误率

那么,过滤器越大越好吗?也不尽然:尽管布隆过滤器在存储方面非常高效(因为它不存储元素,只存储元素是否在集合中),但过滤器大小是预分配的。如果你将过滤器大小设置为 255,且每个数据块至少包含一条带有过滤值的消息,那么每个数据块头部都会分配 255 字节。如果数据块包含许多大消息,这没问题,因为过滤器大小与数据块大小相比微不足道。但如果出现退化情况,比如单条 10 字节消息的数据块配合 10 字节的过滤值,最终得到的过滤器反而比实际数据还要大。

你需要通过自己的用例进行实验,以估算过滤器大小对流大小的影响。流过滤的第一篇文章提供了一个技巧,利用 Stream PerfTest 来估算流的大小(读取整个流而不进行过滤,并查看 rabbitmq_stream_read_bytes_total 指标)。

附赠:AMQP 上的流过滤

尽管访问流的首选方式是流协议,但也支持其他协议,例如 AMQP。任何 AMQP 客户端库也都支持流过滤。

  • 声明:在声明流时,将 x-queue-type 参数设置为 stream,并使用 x-stream-filter-size-bytes 设置过滤器大小。
  • 发布:使用 x-stream-filter-value 头部为外发消息设置过滤值。
  • 消费:使用 x-stream-filter 消费者参数设置期望的过滤值(字符串或字符串数组),并可选地使用 x-stream-match-unfiltered 消费者参数来接收没有任何过滤值的消息(默认为 false)。客户端过滤仍然是必要的!

总结

本文深入介绍了 RabbitMQ 3.13 中的流过滤。它是对第一篇文章的补充,那篇文章介绍了流过滤的使用方法和演示。

流过滤易于使用且收益明显,但了解一些内部机制有助于优化其使用,特别是在复杂的用例中。记住,客户端过滤是必须的,且必须与配置的过滤值保持一致。这通常很容易实现。针对特定用例以最合适的方式设置过滤器大小也是可能的。

© . This site is unofficial and not affiliated with VMware.