跳至主内容
版本:4.3

Stream 插件

概述

流(Stream)是一种持久化且具有副本机制的数据结构,它模拟了具有非破坏性消费者语义的仅追加日志。

此功能在所有当前维护的发布系列中均可用。

流既可以作为常规 AMQP 0.9.1 队列使用,也可以通过专用二进制协议插件及关联的客户端使用。请参阅 RabbitMQ 核心与流插件对比页面了解功能矩阵。

本页面介绍流插件,它允许使用新二进制协议与流进行交互。有关概念概述及流的操作方法,请参阅 RabbitMQ 流指南

另一篇配套指南《流客户端连接》解释了流协议客户端应如何连接到集群节点,以获得最佳的数据局部性和效率(吞吐量、延迟)。

信息

RabbitMQ 流协议的客户端库适用于多种流行平台。其中部分列出如下。标有对勾 (✓) 的为 RabbitMQ 核心团队官方支持的客户端。

使用 流性能测试工具 (Stream PerfTest) 来模拟工作负载并测量 RabbitMQ 流系统的性能。

启用插件

流插件包含在 RabbitMQ 发行版中。在客户端成功连接之前,必须使用 rabbitmq-plugins 启用它。

rabbitmq-plugins enable rabbitmq_stream

插件配置

TCP 监听器

如果不指定配置,流适配器将监听所有接口的 5552 端口,默认用户名/密码为 guest/guest

流监听器所使用的端口可以通过 rabbitmq.conf 进行更改。

以下是一个将监听端口更改为 12345 的最小化配置文件。

stream.listeners.tcp.1 = 12345

而一个仅监听 localhost(IPv4 和 IPv6)的配置将如下所示

stream.listeners.tcp.1 = 127.0.0.1:5552
stream.listeners.tcp.2 = ::1:5552

TCP 监听器选项

该插件支持 TCP 监听器选项配置。

这些设置使用通用前缀 stream.tcp_listen_options,用于控制 TCP 缓冲区大小、入站 TCP 连接队列长度、是否启用 TCP keepalives 等。详情请参阅 网络指南

stream.listeners.tcp.1 = 127.0.0.1:5552
stream.listeners.tcp.2 = ::1:5552

stream.tcp_listen_options.backlog = 4096
stream.tcp_listen_options.recbuf = 131072
stream.tcp_listen_options.sndbuf = 131072

stream.tcp_listen_options.keepalive = true
stream.tcp_listen_options.nodelay = true

stream.tcp_listen_options.exit_on_close = true
stream.tcp_listen_options.send_timeout = 120

心跳超时

心跳超时值定义了 RabbitMQ 和客户端库在多长时间未收到消息后应将对端 TCP 连接视为不可达(宕机)。

RabbitMQ 支持的消息协议使用了类似的机制

流协议连接的默认超时值为 60 秒。

# use a lower heartbeat timeout value
stream.heartbeat = 20

将心跳超时值设置得过低可能会导致误报(即在实际上对端并未故障时,因瞬时网络拥塞、短时间服务器流控等原因将其视为不可用)。

在选择超时值时应考虑到这一点。

根据用户和客户端库维护者多年的反馈,低于 5 秒的值极易导致误报,而 1 秒或更低的值几乎肯定会导致误报。在大多数环境中,5 到 20 秒之间的值是最佳选择。

流控

如果代理无法跟上入站消息的写入和复制速度,快速的发布者可能会压垮代理。因此,每个连接在被阻塞之前允许存在最大数量的未确认消息(initial_credits,默认为 50,000)。当指定数量的消息得到确认(credits_required_for_unblocking,默认为 12,500)时,连接将解除阻塞。您可以根据自己的工作负载更改这些值。

stream.initial_credits = 100000
stream.credits_required_for_unblocking = 25000

这些设置的高值可以提高发布吞吐量,但会增加内存消耗(可能导致代理崩溃)。低值有助于应对大量中速发布的连接。

此设置仅适用于发布者,不适用于消费者。

消费者信用流控

本节介绍流协议的信用流控机制,它允许消费者控制代理分发消息的方式。

消费者在创建其订阅时提供初始信用额度。一个信用点代表代理被允许发送给消费者的一块(chunk)消息。

块(Chunk)是一批消息。这是 RabbitMQ Stream 中使用的存储和传输单位,即消息在块中连续存储,并作为块的一部分进行投递。根据输入情况,一个块可能包含一到几千条消息。

因此,如果消费者创建了一个具有 5 个初始信用点的订阅,代理将发送 5 块消息。代理每投递一个块,就会减去一个信用点。当订阅的信用点用完时,代理将停止发送消息。在我们的例子中,代理在投递 5 块消息后将停止发送。这通常不是我们想要的,因此消费者可以为其订阅提供信用点以获取更多消息。

根据处理消息的速度,由消费者(即客户端库和/或应用程序)来提供信用点。我们希望消息能够持续流动,因此一个经验法则是:创建订阅时至少提供 2 个信用点,并在收到每一块新消息时补充一个信用点。这样,网络上应该始终有消息在流动,消费者应该始终处于忙碌状态,而不是空闲。

通过这种信用流控机制,消费者可以选择代理如何向其投递消息。这有助于避免消费者过载或闲置。消费者信用流控如何向应用程序公开取决于客户端库,没有服务器端的设置可以改变其行为。

通告的主机和端口

流协议允许发现流的拓扑结构,即给定一组流的领导者和副本在集群中的位置。这样客户端可以选择连接到适当的节点与流进行交互:连接领导者节点以发布,连接副本节点以消费。默认情况下,节点会返回其主机名和监听端口,这在大多数情况下适用,但在某些情况下(如集群节点与客户端之间的代理、运行在容器中的集群节点和/或客户端等)可能不适用。

advertised_hostadvertised_port 键允许指定代理节点在被询问流拓扑时返回哪些信息。可以根据基础设施设置这些配置,以便客户端能够连接到集群节点。

stream.advertised_host = rabbitmq-1
stream.advertised_port = 12345

流客户端连接指南涵盖了为什么在某些部署中需要 advertised_hostadvertised_port 设置。

这些配置条目具有对应的 TLS 设置advertised_tls_hostadvertised_tls_port)。

最大帧大小

RabbitMQ 流协议使用最大帧大小限制。默认值为 1 MiB,如有必要,可以增加该值。

# in bytes
stream.frame_max = 2097152

TLS 支持

要将 TLS 用于流连接,必须在代理中配置 TLS。要启用支持 TLS 的流连接,请使用 stream.listeners.ssl.* 配置键添加流的 TLS 监听器。

该插件将使用核心 RabbitMQ 服务器证书和密钥(就像 AMQP 0.9.1 和 AMQP 1.0 监听器一样)

ssl_options.cacertfile = /path/to/tls/ca_certificate.pem
ssl_options.certfile = /path/to/tls/server_certificate.pem
ssl_options.keyfile = /path/to/tls/server_key.pem
ssl_options.verify = verify_peer
ssl_options.fail_if_no_peer_cert = true

stream.listeners.tcp.1 = 5552
# default TLS-enabled port for stream connections
stream.listeners.ssl.1 = 5551

此配置创建了一个 5552 端口的标准 TCP 监听器和一个 5551 端口的 TLS 监听器。

当设置了 TLS 监听器时,您可能希望停用所有非 TLS 监听器。配置如下:

stream.listeners.tcp = none
stream.listeners.ssl.1 = 5551

普通连接一样,可以配置通告的 TLS 主机和端口。使用 TLS 时,插件返回以下元数据:

  • hostname:如果已设置,则为 advertised_host;如果未设置 advertised_host,则为真实主机名。
  • port:当前的 TLS 端口。

可以通过同时或单独设置 advertised_tls_hostadvertised_tls_port 配置条目来覆盖此行为。

stream.advertised_host = private-rabbitmq-1
stream.advertised_port = 12345
stream.advertised_tls_host = public-rabbitmq-1
stream.advertised_tls_port = 12344
© . This site is unofficial and not affiliated with VMware.