【问题标题】:Queue Fairness and Messaging Servers队列公平和消息服务器
【发布时间】:2016-09-20 18:50:41
【问题描述】:

我正在寻找解决消息服务器和队列的 FIFO 特性的问题。在某些情况下,我希望将队列中的消息分配给消费者池,而不是按照消息传递顺序的标准。理想情况下,这将防止用户占用系统中的共享资源。以这个过于简化的场景为例:

  • 应用程序中有一个功能,用户可以清空他们的垃圾桶。
  • 此事件为垃圾箱中的每个项目分派 DELETE 消息
  • 此队列的消费者调用具有速率限制 API 的 Web 服务。

鉴于每个用户的垃圾箱中都可能有大量消息,我们有哪些选项可以允许同时处理每个垃圾箱而不考虑排队时间?在我看来,有一些明显的解决方案:

  • 为每个用户创建一个单独的队列和消费者池
  • 将消息从单个队列随机传递到单个消费者池

在我们的例子中,为每个用户创建一个单独的队列并管理消费者确实不切实际。可以做到,但我认为如果合理的话,我真的更喜欢第二种选择。我们正在使用 RabbitMQ,但如果有更适合此任务的技术,则不一定要与之绑定。

我正在考虑使用 Rabbit 的消息优先级来帮助随机发送的想法。通过随机为消息分配 1 到 10 之间的优先级,这应该有助于分发消息。这种方法的问题是,如果队列从未完全清空,则优先级最低的消息可能会永远卡在队列中。我以为我可以在消息上使用 TTL,然后以升级的优先级重新排队消息,但我在 docs 中注意到了这一点:

应该过期的消息仍然只会从头部过期 队列。这意味着与普通队列不同,即使是每个队列 TTL 可能导致过期的低优先级消息被卡在后面 未过期的更高优先级的。这些消息永远不会 已交付,但它们会出现在队列统计信息中。

我担心我可能会用这种方法进入兔子洞。我想知道其他人是如何解决这个问题的。任何关于创意路由、消息传递模式或任何替代解决方案的反馈都将受到重视。

【问题讨论】:

    标签: rabbitmq jms message-queue messaging


    【解决方案1】:

    一种解决方案是插入Resequencer。该链接中的诊断概述了该原理。在您的情况下,类似于:

    • 应用程序将其 DELETE 消息按原样分派到删除队列中。
    • Resequencer(您编写的一个新组件)介于原始发布者和原始消费者之间。它:

      • 将消息从 DELETE 队列拉到内存中
      • 按用户将它们放入(内存中)队列中
      • 将它们重新发布到新队列(例如 FairPriorityDeleteQueue),循环以公平地交错来自不同原始用户的任何消息
      • 限制其在 FairPriorityDeleteQueue 中的重新发布速率,以使 FairPriorityDeleteQueue 的长度(可通过定期轮询 rabbitmq 管理 api 获得)永远不会超过您选择的某个整数 N,或者限制为与消费者的速率限制删除 API 相关的某个速率采用。
      • 不确认从原始 DELETE 队列中提取的任何消息,直到将其重新发布到 FairPriorityDeleteQueue(因此您永远不会丢失消息)
    • 原来的消费者订阅了 FairPriorityDeleteQueue。

      • 您将这些消费者的 preFetchCount 设置得相当低 (

    --

    需要注意的几点:

    • 发布到 FairPriorityDeleteQueue 和/或从 FairPriorityDeleteQueue 中提取消息的速率或长度限制是必不可少的。如果您不进行限制,Resequencer 可能会在收到消息后尽快发送消息,从而限制重新排序的可能性。
    • Resequencer 在重新排序时当然充当一种内存缓冲区。如果原始发布者可以突然将 非常 大量消息发布到队列中,您可能需要对 Resequencer 进程进行内存限制,使其不会摄取超过其可以容纳的内容。

    您有一个限制吞吐量的外部因素(最终删除 API)这一事实极大地帮助了您的特定场景。 没有这样一个外部限制因素,选择最佳参数会困难得多对于这样的重新排序器,可以在特定环境中平衡吞吐量与重新排序。

    【讨论】:

    • 最终,我得到了一个解决方案,它最终通过负载平衡器对来自不同队列的消息重新排序,但它避免了额外的队列和内存中的重新排序。我会尽快添加答案。
    • @MikeCantrell 对。同意这避免了额外的队列。从文章的快速浏览来看,代码使用 过去 队列内容作为将消息从队列中拉出的优先级指南? (是吗?)
    【解决方案2】:

    所以我最终从网络路由器手册中取出一页。这是他们路由器需要解决的一个问题,以允许公平的流量模式。 This video 很好地分解了问题和解决方案。

    将问题翻译到我的领域:

    以及解决方案:

    负载平衡器是通道和已知数量的队列的包装器,它使用加权算法在每个队列上接收的消息之间进行平衡。我们找到了一个really interesting article/implementation,目前看来运行良好。

    使用此解决方案,我还可以在发布消息后对工作区进行优先级排序,以提高其吞吐量。这是一个非常好的功能。

    摆在我面前的最大挑战是队列管理。将有太多的队列离开绑定到交换很长一段时间。我正在开发一些工具来管理它们的生命周期。

    【讨论】:

    • 太棒了!做得好,从网络适应软件消息传递。 :) 下次遇到类似情况时,我将不得不记住这一点
    • 这是一个非常好的答案,赞!
    【解决方案3】:

    我认为在这种情况下不需要重新排序器。如果您需要确保按特定顺序删除项目,也许是这样。但这只有在您大致同时发送多条消息并且需要保证消费者端的顺序时才会发挥作用。

    由于您提到的原因,您还应该避免超时情况。 timeout 是为了告诉 RabbitMQ 一条消息不需要被处理——或者它需要被路由到一个死信队列,以便我可以被其他一些代码处理。虽然您可以让超时工作,但我认为这不是一个好的选择。

    优先级可能会解决部分问题,但可能会引入文件永远不会被处理的情况。如果您将优先级 1 的消息放回队列中的某个位置,并且您继续将优先级 2、3、5、10 等放入队列,则可能不会处理 1。正如您所指出的,超时并不能解决这个问题。

    为了我的钱,我会建议一种不同的方法:为单个文件串行发送删除请求。

    即发送1条消息删除1个文件。等待回复说它已经完成。然后发送下一条消息以删除下一个文件。

    这就是我认为这会起作用的原因,以及如何管理它:

    长时间运行的工作流程,单个文件删除请求

    在这种情况下,我建议使用“传奇”(又名a long-running workflow object)的想法采取多步骤方法来解决问题。

    当用户请求删除他们的垃圾箱时,您通过 rabbitmq 向可以处理删除过程的服务发送一条消息。该服务会为该用户的垃圾桶创建一个 saga 实例。

    传奇收集了垃圾箱中需要删除的所有文件的列表。然后它开始发送删除单个文件的请求,一次一个。

    对于删除单个文件的每个请求,saga 都会等待响应说文件已被删除。

    当 saga 收到上一个文件已被删除的消息时,它会发出下一个删除下一个文件的请求。

    一旦所有文件都被删除,saga 会更新自己和系统的任何其他部分,说垃圾桶是空的。

    处理多个用户

    当您有一个用户请求删除时,对他们来说事情会很快发生。他们很快就会清空垃圾。

    u1 = 用户 1 垃圾箱删除请求 |u1|u1|u1|u1|u1|u1|u1|u1|u1|u1完成|

    当您有多个用户请求删除时,一次发送一个文件删除请求的过程意味着每个用户将有相同的机会获得下一个文件删除。

    u1 = 用户 1 垃圾箱删除请求 u2 = 用户 2 垃圾箱删除请求 |u1|u2|u1|u1|u2|u2|u1|u2|u1|u2|u2|u1|u1|u1|u2|u2|u1|u2|u1|u1done|u2|u2done|

    这样,将共享使用资源来删除文件。总体而言,清空每个人的垃圾桶需要更长的时间,但他们会更快看到进展,这是人们认为系统快速/响应他们的请求的一个重要方面。

    优化小文件集与大文件集

    在用户数量较少且文件数量较少的情况下,上述解决方案可能会比一次性删除所有文件要慢。毕竟,通过rabbitmq发送的消息会更多——每个需要删除的文件至少有2条消息(一个删除请求,一个删除确认响应)

    要进一步优化,您可以做几件事:

    1. 在像这样拆分工作之前,有一个最小的垃圾桶大小。低于最低限度,您只需一次将其全部删除

    2. 将工作分成几组文件,而不是一次一个。也许 10 或 100 个文件会比一次 1 个文件更好的组大小

    这些解决方案中的任何一个(或两个)都可以通过减少发送的消息数量和稍微分批处理工作来帮助提高流程的整体性能。

    您需要在真实场景中进行一些测试,以了解其中哪些(或可能两者兼有)会有所帮助以及在什么设置下会有所帮助。

    多用户问题

    还有一个您可能面临的额外问题 - 许多用户。如果您有 2 或 3 个用户请求删除,这没什么大不了的。

    但是,如果您有 100 或 1000 个用户请求删除,则个人可能需要很长时间才能清空他们的垃圾桶。

    对于这种情况,您可能需要一个更高级别的控制流程,所有清空垃圾桶的请求都将由另一个 Saga 管理。这个 saga 会限制活跃的垃圾桶删除 sagas 的数量。

    例如,如果您有 10 个删除垃圾箱的活动请求,则限速 saga 只会启动其中的 3 个,并且会等待其中一个完成,然后再开始下一个。

    同样,出于性能原因,您需要测试您的实际场景以查看是否需要这样做,并查看应该有哪些限制。


    在您的实际场景中可能需要考虑其他场景,但我希望这能让您走上正轨! :)

    【讨论】:

    • 不幸的是,除非我同时运行它们,否则我将无法为每个垃圾桶实现不错的删除率。例如,以每条消息 50 毫秒的速度删除包含 10k 条消息的垃圾箱需要 8 分钟以上。在这种情况下,批处理可能会有所帮助,但我的情况相同,每个工作单元都需要在不同的服务器上完成(由于一些繁重的文件 IO 和 CPU 匹配)。
    • 在不同的系统上完成工作很容易。每个消费者设置消费者预取限制为 1,并且在工作完成之前不确认消息。但是我可以看到,如果您要批量处理 10 个项目并且每个项目都需要在单独的服务器上运行,这仍然会有问题。
    猜你喜欢
    • 1970-01-01
    • 2014-02-23
    • 2011-05-23
    • 1970-01-01
    • 2012-11-18
    • 2013-06-06
    • 2010-11-02
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多