使用 RabbitMQ 防止无界缓冲区
我们架构中的不同服务在运行过程中需要一定量的资源,无论是 CPU、RAM 还是磁盘空间,我们都需要确保有足够的资源。如果我们不对服务器将要使用的资源数量设置限制,迟早会遇到麻烦。当数据库耗尽文件系统空间、媒体存储填满图片且从不转移它们,或者 JVM 耗尽 RAM 时,就会发生这种情况。即使是备份解决方案,如果不设置过期/删除旧备份的策略,也会成为一个问题。嗯,队列也不例外。我们必须确保我们的应用程序不会让队列无休止地增长。我们需要有某种策略来删除/驱逐/迁移旧消息。
为什么会出现这个问题?
队列中消息堆积的原因有很多。首要原因是数据生产者的速度超过了消费者的处理速度。幸运的是,解决方案很简单:增加更多的消费者。
如果我们的应用程序仍然无法处理负载该怎么办?例如,您的消费者处理每条消息的时间太长,而且您无法增加更多的消费者(因为服务器资源耗尽了)。那么,队列将开始堆积消息。RabbitMQ 针对快速消息传递进行了优化,其队列中保留的消息越少越好。虽然 RabbitMQ 自带了多种流控机制,但您当然还是希望有一种方法来防止触发流控。让我们看看 RabbitMQ 如何在这方面提供帮助。
队列消息 TTL(Per-Queue Message TTL)
RabbitMQ 允许设置队列级别的消息 TTL(生存时间),这会使服务器不再投递在队列中停留时间超过设定 TTL 的消息。此外,服务器将尝试尽快过期或死信这些消息。
当您的数据只有在及时到达时才对生产者有意义时,这种方法非常有效。如果您的数据不能丢弃,但您仍希望队列尽可能保持为空,请参阅下文关于“死信(Dead lettering)”的章节。
有两种设置队列 TTL 的方法,一种是在 queue.declare 期间传递一些额外的参数,如下所示:
Map<String, Object> args = new HashMap<String, Object>();
args.put("x-message-ttl", 60000);
channel.queueDeclare("myqueue", false, false, false, args);
上述代码将告诉 RabbitMQ 在 60 秒后让队列 myqueue 上的消息过期。
也可以通过为队列添加策略来设置相同的内容:
rabbitmqctl set_policy TTL ".*" '{"message-ttl":60000}' --apply-to queues
此策略将匹配默认虚拟主机中的所有队列,并使消息在 60 秒后过期。请注意,Windows 命令略有不同。当然,您可以让该策略仅匹配一个队列。更多详细信息请见:参数和策略。
如果我们想要更精细地控制哪些消息过期该怎么办?
消息 TTL(Per-Message TTL)
RabbitMQ 还支持设置单条消息的 TTL。我们可以通过在 basic.publish 方法调用中设置 expiration 字段来设置消息的 TTL。与前面的情况一样,该值应以毫秒为单位。以下代码将发布一条在 60 秒后过期的消息:
byte[] messageBodyBytes = "Hello, world!".getBytes();
AMQP.BasicProperties properties = new AMQP.BasicProperties();
properties.setExpiration("60000");
channel.basicPublish("my-exchange", "routing-key", properties, messageBodyBytes);
如果我们结合使用消息 TTL 和队列 TTL,则以较短的 TTL 为准。RabbitMQ 将确保消费者永远不会接收到过期的消息;但在使用消息 TTL 的情况下,在这些消息到达队列头部之前,它们不会被过期处理。
队列 TTL(Queue TTL)
使用 RabbitMQ,我们还可以让整个队列过期,即在队列处于未使用状态一段特定时间后将其删除。假设我们将队列设置为一小时后过期。如果在这一个小时内,队列上没有消费者、没有发出 basic.get 命令,或者队列没有被重新声明,那么 RabbitMQ 会将其视为未使用并将其删除。
如果您在用户在线时为他们创建队列,但希望在 15 分钟不活动后删除这些队列,那么您可能希望使用此功能。考虑一个为每个已连接用户保留一个队列的聊天应用程序。您可以声明一个 auto_delete 队列,它会在用户关闭通道后立即消失;但这在某些场景下可能有用,但如果用户因为移动网络连接质量差而断开连接,会发生什么?您肯定不希望在他们断开连接后立即删除他们所有的消息。有了这个功能,您可以让这些队列存活更长的时间。
以下是使用 Java 客户端设置 15 分钟队列过期时间的方法:
Map<String, Object> args = new HashMap<String, Object>();
args.put("x-expires", 900000);
channel.queueDeclare("myqueue", false, false, false, args);
也可以通过策略设置:
rabbitmqctl set_policy expiry ".*" '{"expires":900000}' --apply-to queues
队列长度限制
如果我们希望队列中的消息不超过某个阈值,我们可以通过在声明队列时使用 x-max-length 参数进行配置。这是一种控制容量的简洁且简单的方法;如果队列达到阈值且有新消息到达,则位于队列头部(即“较旧的消息”)将被丢弃,从而为新消息腾出空间。这种行为的原因之一是旧消息对于您的应用程序可能已经不重要了,所以新消息被允许进入队列。
请记住,队列长度只计算已准备好投递的消息。未确认(Unacked)的消息不会计入总数。拥有正确的 basic.qos 设置将对您的应用程序有所帮助,因为默认情况下 RabbitMQ 会尽可能多地向消费者发送消息,这会导致一种情况:您的队列看起来是空的,但实际上您有很多未确认的消息,它们也在占用资源。
设置队列长度限制非常容易,以下是一个在 Java 中设置 10 条消息限制的示例:
Map<String, Object> args = new HashMap<String, Object>();
args.put("x-max-length", 10);
channel.queueDeclare("myqueue", false, false, false, args);
也可以通过策略设置:
rabbitmqctl set_policy Ten ".*" '{"max-length":10}' --apply-to queues
混合使用策略
请记住,在任何给定时间,最多只有一个策略应用于一个队列。因此,如果您连续运行之前的 set_policy 命令,只有最后一个会生效。将多个策略应用于同一资源的技巧在于将所有策略一起放在同一个 JSON 对象中,例如:
rabbitmqctl set_policy capped_queues "^capped\." \
'{"max-length":10, "expires":900000, "message-ttl":60000}' --apply-to queues
完全不排队
等等,我读对了吗?是的。不排队。
想象一下,在一个非常忙碌的日子里,您来到邮局,发现每个柜台都在忙。由于您没有时间浪费在排队上,您就直接离开并继续做您之前做的事情。换句话说:您有一个必须立即被服务的请求,即:不进行排队。好吧,RabbitMQ 可以通过您的应用程序消息和队列实现类似的效果。
诀窍在于将 per-queue-TTL 设置为 0(零)。如果消息无法立即投递给消费者,那么它们将立即过期。如果您设置了死信交换机(dead-letter exchange),那么您可以将这些消息发送到单独的队列中。
死信(Dead Lettering)
我们已经多次提到过死信。此功能的作用是,您可以为其中一个队列设置死信交换机 (DLX),然后当该队列中的消息过期或超出队列限制时,消息将被发布到 DLX。您可以自行绑定一个单独的队列到该交换机,然后稍后处理发送到那里的消息。
以下是用于设置 DLX 的 queue.declare 示例:
channel.exchangeDeclare("some.exchange.name", "direct");
Map<String, Object> args = new HashMap<String, Object>();
args.put("x-dead-letter-exchange", "some.exchange.name");
channel.queueDeclare("myqueue", false, false, false, args);
死信消息将使您的队列保持适当的大小和预期的消息量,但这并不能防止节点被消息填满。如果这些消息被排队到同一节点上的不同队列中,那么在某个时刻,这个新的死信队列可能会引发问题。在这种情况下,您可以做的是使用交换机联邦(exchange federation)将这些消息发送到单独的节点,并将其与应用程序的主要流程分开处理。
结论
关于请求到达系统的排队理论,其中一个基本问题可以表述如下1
λ = 平均到达时间
µ = 平均服务速率
如果 λ > µ 会发生什么?
队列长度随时间推移趋向于无穷大。
我们知道,如果我们的架构中在任何一点遇到这个问题,迟早我们的应用程序会陷入困境。幸运的是,RabbitMQ 提供了许多特性,如队列和消息 TTL、队列过期和队列长度限制,专门用于避免此问题。更有趣的是,我们不需要仅仅因为使用这些功能就丢失消息。死信交换机可以帮助我们将消息重定向到更合适的地方。现在是我们把这些技术加入到我们的队列和消息传递工具库中的时候了。