【发布时间】:2021-05-18 03:24:15
【问题描述】:
我正在实现一个消费者,它处理来自消息顺序很重要的队列中的消息。我想使用 NodeJS 实现一种机制,其中:
- 消费者函数正在消费队列中的消息
m1, m2, ..., mN - 执行 IO 密集型操作并处理消息。
m -> m' - 将结果
m'存储在redis 缓存中。 - 在每个消息处理后确认队列 (2)
在另一个函数中,我正在从缓存中收听消息
- 将处理后的消息
m'发送到外部系统 - 如果外部系统能够处理外部系统,则从缓存中删除处理后的消息
- 如果外部系统拒绝已处理的消息,则停止发送消息,丢弃缓存中未发送的已处理消息并将偏移量重置为队列中最后接受的消息。例如,如果
m12'是系统接受的最后一条消息,并且我已经从队列中确认了m23,那么我必须将m13'丢弃到m23'并重置偏移量,以便消费者可以读取并启动再次从m13处理。
几个假设:
-
m到m'的处理非常密集,我很乐观地处理它们,因为我知道大多数时候不会出现故障
根据当前的假设和目标,我有什么方法可以使用 RabbitMQ 或任何 Azure 等效项来实现这一目标?我的客户不喜欢 Kafka 或任何 Azure 等效的 Kafka(Azure 事件中心)。
【问题讨论】:
-
尝试强制执行订单或处理以及像 Azure 事件中心这样的大规模摄取的主要问题是,这样做是有代价的。如果在发送方失败并重试的情况下节点接受请求的速度太慢,是否应该放弃尝试,或者是否应该将其发送到另一个节点,在该节点可能在前一条消息阻止所有内容之前对其进行处理?如果您可以设计一个具有弹性的发送方,因为它可以忽略失败的发送并重试,那么就有可能在集线器级别强制执行处理顺序,否则您将需要重新排序消息
-
我对此想得越多,你的解释可能已经抽象得太多了,让两个消费者以不同的速率处理并试图告诉一个消费者重置回以前的状态意味着你将不得不将以前的数据保留很长时间。
-
这个场景有多少个消息生产者?
-
非常感谢克里斯考虑这个问题。可以有多个消息生产者。为简单起见,我认为只有一个生产者。因为我不能在这里进行负载平衡,因为顺序很重要。对于持久化数据的问题,不会超过一天左右。我期待每秒数百条消息。我曾考虑将订单存储在单独的数据库中并使用它进行重播,但这对我来说似乎不是一个优雅的解决方案。