【问题标题】:Replaying Messages in Order按顺序重播消息
【发布时间】:2021-05-18 03:24:15
【问题描述】:

我正在实现一个消费者,它处理来自消息顺序很重要的队列中的消息。我想使用 NodeJS 实现一种机制,其中:

  1. 消费者函数正在消费队列中的消息m1, m2, ..., mN
  2. 执行 IO 密集型操作并处理消息。 m -> m'
  3. 将结果m' 存储在redis 缓存中。
  4. 在每个消息处理后确认队列 (2)

在另一个函数中,我正在从缓存中收听消息

  • 将处理后的消息m' 发送到外部系统
  • 如果外部系统能够处理外部系统,则从缓存中删除处理后的消息
  • 如果外部系统拒绝已处理的消息,则停止发送消息,丢弃缓存中未发送的已处理消息并将偏移量重置为队列中最后接受的消息。例如,如果m12' 是系统接受的最后一条消息,并且我已经从队列中确认了m23,那么我必须将m13' 丢弃到m23' 并重置偏移量,以便消费者可以读取并启动再次从m13 处理。

几个假设:

  • mm' 的处理非常密集,我很乐观地处理它们,因为我知道大多数时候不会出现故障

根据当前的假设和目标,我有什么方法可以使用 RabbitMQ 或任何 Azure 等效项来实现这一目标?我的客户不喜欢 Kafka 或任何 Azure 等效的 Kafka(Azure 事件中心)。

【问题讨论】:

  • 尝试强制执行订单或处理以及像 Azure 事件中心这样的大规模摄取的主要问题是,这样做是有代价的。如果在发送方失败并重试的情况下节点接受请求的速度太慢,是否应该放弃尝试,或者是否应该将其发送到另一个节点,在该节点可能在前一条消息阻止所有内容之前对其进行处理?如果您可以设计一个具有弹性的发送方,因为它可以忽略失败的发送并重试,那么就有可能在集线器级别强制执行处理顺序,否则您将需要重新排序消息
  • 我对此想得越多,你的解释可能已经抽象得太多了,让两个消费者以不同的速率处理并试图告诉一个消费者重置回以前的状态意味着你将不得不将以前的数据保留很长时间。
  • 这个场景有多少个消息生产者?
  • 非常感谢克里斯考虑这个问题。可以有多个消息生产者。为简单起见,我认为只有一个生产者。因为我不能在这里进行负载平衡,因为顺序很重要。对于持久化数据的问题,不会超过一天左右。我期待每秒数百条消息。我曾考虑将订单存储在单独的数据库中并使用它进行重播,但这对我来说似乎不是一个优雅的解决方案。

标签: node.js azure rabbitmq


【解决方案1】:

在消息将总是按顺序生成的情况下,您可能只需要一个简单的队列。

Azure Queues 很容易进入,但队列的一般操作模式是在消息成功处理后将其删除。

如果您可以避免必须“回滚”或从较早时间重新处理的情况,那么如果您可以避免编排方面,那么这将是一个更简单的选择。

这是您将难以解决的“回去重播”。如果您可以按顺序模式实现两个队列,其中处理来自一个队列的消息成功地将消息推送到下一个队列,那么我们永远不需要返回,因为辅助消费者永远无法在主消费者之前处理。


使用 Azure 事件中心,重置偏移以进行处理要容易得多,因为无论消息的 read 状态如何,消息都会保留在存储桶中,(实际上任何给定的消息没有这样的状态)并且消费者自己维护偏移量指针。它还支持多个消费者组,这将使每个消费者都可以使用消息的副本。

您可以制定计划,在不超出预算的情况下将数据保留长达 7 天。

对于您的用例,Azure 事件中心等大规模遥测摄取服务存在两个问题

  1. 对于非常接近的消息,消息的接收顺序不太可靠,Hub 旨在同时接收来自多个来源的许多消息,因此其内部架构不太关心尝试保持精确的顺序,它在消息上记录了准确的接收时间戳,但它不保证整个记录序列将与您要按接收时间戳排序的场景完全匹配。 (这是一个微妙但重要的区别)

  2. 事件中心(以及许多客户端处理代码示例)旨在保证跨多个并发消费线程的 Exactly Once 交付。再次鼓励消费者是异步的,服务将尝试确保下一个可用线程重试失败的处理尝试。

因此您可以使用事件中心,但您必须绕过或禁用它的许多功能,这通常是一个强烈的信息,表明它不适合您的目的,但如果您想探索它,您会想要限制并发方面:

  • 最小化分区数

    您可能希望为每个消息生产者使用 1 个分区,或者至少为每个顺序集设置一个分区,在单个分区内维护序列更简单

  • 确保您的消息发送者(生产者)仅发送到特定分区

    每个生产者必须使用唯一的分区键

  • 为您的每个消费者创建一个消费者组
  • 一次处理一条消息,而不是批量处理
  • 单线程处理

我在为工业物联网(PLC 遥测)和农业物联网 (Raspberry Pi) 设备实施设计基于 MS Azure 的解决方案方面拥有丰富的经验。在几乎所有情况下,我们认为消息的顺序很重要,但除非您保持实时的 2 路命令和控制,否则您通常可以采用一种乐观的方法,其中每条消息和任何衍生品在传输时正确或正确。

如果设备在任何时间段内都存在离线的可能性,那么在设备恢复在线时处理通过系统刷新的陈旧数据确实可以对顺序逻辑编程造成严重影响。

退后一步来分析您的解决方案,EventHubs 确实提供了一种方便的方式来将处理回滚到先前的偏移量,只要该记录仍在 bucket 中,但您可以重新 -设计您的逻辑流程,以便您不必重新处理旧数据?

驱动这个序列的要求是什么?如果保持顺序非常重要,那么您可能应该使用一个可以处理所有事情的消费者来处理数据,或者考虑以顺序方式链接队列。

【讨论】:

    猜你喜欢
    • 2021-02-05
    • 2014-08-23
    • 2018-11-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多