【问题标题】:Reliable Webhook dispatching system可靠的Webhook调度系统
【发布时间】:2021-08-21 14:59:05
【问题描述】:

我很难为 webhook 调度系统找到一个可靠且可扩展的解决方案。

当前系统使用 RabbitMQ 和 webhook 队列(我们称之为 events),它们被消费和分派。该系统运行了一段时间,但现在出现了一些问题:

  • 如果系统用户产生的事件过多,会占用队列,导致其他用户长时间收不到 webhook
  • 如果我将所有事件分成多个队列(通过 URL 哈希),它会降低出现第一个问题的可能性,但是当非常忙碌的用户访问同一个队列时,它仍然会不时发生
  • 如果我尝试将每个 URL 放入自己的队列中,挑战是动态地创建/分配消费者到这些队列。就RabbitMQ 文档而言,API 在过滤非空队列或未分配消费者的队列方面非常有限。
  • Kafka 而言,据我了解,通过阅读有关它的所有内容,情况在单个分区的范围内将是相同的。

所以,问题是 - 有没有更好的方法/系统来达到这个目的?也许我错过了一个非常简单的解决方案,允许一个用户不干扰另一个用户?

提前致谢!

【问题讨论】:

  • 我觉得散列是正确的解决方案。您可以实施传入速率限制,以防止不良行为者减慢特定队列/分区的速度
  • 传入的速率限制不会减慢生产者的速度吗?此外,这意味着“慢”消息无论如何都需要转到其他地方。
  • 我不明白如何使用 url hasing 将事件拆分为多个队列。请给个解释好吗?
  • @nsv 每个 webhook 处理程序都有一个唯一的 URL。每个 webhook 处理程序都可以分配多个事件。因此,当一个事件被创建时,它会被放入其各自的 webhook 处理程序的队列中,并且由于每个 webhook 处理程序都有一个唯一的 URL,所以它基本上是相同的。
  • @Arthur 但是当你有这么多网址时,如果每个网址都有一个队列,你如何管理?

标签: apache-kafka rabbitmq webhooks event-dispatching


【解决方案1】:

您可以试验几个 rabbitmq 功能来缓解您的问题(不完全删除它):

  • 使用公共random exchange 将事件拆分到多个队列中。它将缓解大量事件峰值并将工作分派给多个消费者。

  • 为您的队列设置一些TTL policies。这样,如果处理速度不够快,Rabbitmq 可能会将事件重新发布到另一组队列(例如通过另一个私有随机交换)。

您可能有多个事件“周期”,配置不同(即周期数和每个周期的 TTL 值)。您的第一个周期会尽其所能处理新事件,通过随机交换下的多个队列来缓解峰值。如果它未能足够快地处理事件,则将事件移动到具有专用队列和消费者的另一个循环。

这样,您可以确保新事件有更好的更改以快速处理,因为它们始终会在第一个周期中发布(而不是在其他用户的一堆旧事件之后)。

【讨论】:

  • 嗯,是的,但这也会引入单个用户开始并行接收事件的可能性,从而打破“事件顺序”。
  • 确实如此。在阅读您的帖子时,我没有意识到这是一个问题。我想它会给你留下“将每个 URL 放在自己的队列中”的解决方案。
  • 那么主要的挑战是 - 如何动态地将消费者连接到这些队列?以及如何在所有队列之间平均分配所有消费者?
  • 你的消费者难道不能“知道”你想要的队列分配,所以他们可以在启动时创建它们并在工作完成时删除它们吗?
  • 这就是这里的问题:D 有没有办法在动态创建的队列中平均分配消费者?这样每个消费者都会接受一个没有消费者的队列并在其上工作,然后移动到下一个?
【解决方案2】:

如果您需要订购,不幸的是您依赖于用户输入。

但在卡夫卡的世界里,这里有几件事要提一下;

  • 您可以使用Transactions 实现exactly-once 交付,这使您可以构建类似于常规 AMQP 的类似系统。
  • Kafka 支持按键分区。这允许您保持相同键的处理顺序(在您的情况下为 userId)。
  • 可以通过调整所有生产者、服务器和消费者端(批量大小、飞行请求等。有关更多参数,请参阅 Kafka documentation)来增加吞吐量。
  • Kafka 支持消息压缩,这可以减少网络流量并增加吞吐量(对于 LZ4 等快速压缩算法只会消耗更多 CPU 资源)。

在您的场景中,分区是最重要的。您可以增加分区以同时处理更多消息。你的消费者可以和你在同一个消费者组中的分区一样多。即使您在达到分区数后进行扩展,您的新使用者也将无法读取并且他们将保持未分配状态。

与常规的 AMQP 服务不同,Kafka 不会在您阅读消息后删除它,只是标记 consumer-gorup-id 的偏移量。这使您可以同时做几件事。就像在单独的进程中计算实时用户数一样。

【讨论】:

  • 不是我想要的。我通过如何使用 RabbitMQ 处理这个问题的方式找到了解决方案。稍后会发布答案。
  • @Arthur 我对你的解决方案很感兴趣,如果你能找到一些时间在这里分享它。
  • @DavidL 我已经发布了问题的答案
【解决方案3】:

所以,我不确定这是否是解决这个问题的正确方法,但这就是我想出的。

先决条件:带有重复数据删除插件的 RabbitMQ

所以我的解决方案包括:

  • g:events 队列 - 我们称之为 parent 队列。此队列将包含所有需要处理的child 队列的名称。可能它可以用其他一些机制(比如 Redis sorted Set 什么的)来代替,但是你必须自己实现 ack 逻辑。
  • g:events:<url> - 有 child 队列。每个队列只包含需要发送到url 的事件。

将 webhook 有效负载发布到 RabbitMQ 时,您将实际数据发布到 child 队列,然后另外将 child 队列的名称发布到 parent 队列。重复数据删除插件不允许发布相同的child 队列两次,这意味着只有一个消费者可能会收到该child 队列进行处理。

所有消费者都在消费parent队列,收到消息后,他们开始消费消息中指定的child队列。在child 队列为空后,您确认parent 消息并继续。

此方法允许非常精细地控制允许处理哪些child 队列。如果某些child 队列占用了太多时间,只需ack parent 消息并将相同的数据重新发布到parent 队列的末尾。

我知道这可能不是最有效的方法(不断发布到parent 队列也有一些开销),但它就是这样。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-03-07
    • 1970-01-01
    • 2023-04-01
    • 1970-01-01
    • 1970-01-01
    • 2011-10-31
    • 1970-01-01
    • 2012-08-12
    相关资源
    最近更新 更多