如何处理消息队列中无序出现的消息?

Mik*_*ike 6 distributed-computing message-queue distributed-system apache-kafka

我曾经在一次采访中被问到,如何处理消息队列中无序传入的消息。已经有一段时间了,我还没有找到明确的答案,我想知道该领域的专家是否可以帮助我回答这个问题,以满足我自己的好奇心。

据我了解,某些消息队列提供一次性和 FIFO 保证。我还知道流系统中事件时间和处理时间的概念。例如,在像 Kafka 这样的基于日志的消息队列中,由于偏移量和消息持久性的存在,混合排序可能不太可能发生(我可能是错的)。我还考虑过使用时间戳,要求每个消息发送者在发送之前记录消息的时间,但由于时钟偏差,这充满了不一致。

考虑到所有这些,我想知道如何解决 AMQP、JMS 或 RabbitMQ 等传统消息传递系统中的混合排序问题,其中十几个物联网设备可能正在发送消息,而我作为消费者希望以正确的顺序协调它们。

Web*_*ver 4

如果您的系统正在使用队列,提供有序消息保证,那么只需使用该通道(如 kakfa 的单个分区,某些设置下的 AMQP)。但是,如果您的系统使用的队列不提供严格的排序,那么一般的想法是客户端可以在发送到队列的每条消息上附加单调递增的[1]数字(或时间戳)。这构成了生产者打算发送给接收者的序列的基础。

如何获得单调递增的价值:

使用时间戳: 带有 CLOCK_MONOTONIC[2] 的 POSIX Clock_gettime() 函数提供了获取单调递增时间戳的选项,生产者可以使用它在每条消息上添加时间戳。当接收方发现收到的消息的时间戳早于最新消息时,它可以识别出乱序消息。

使用序列号: 在发送每条消息之前,您可以简单地增加一个原子计数器并将计数器值附加到每条消息,以便接收者可以知道预期的顺序。这将形成严格递增的序列。方法与 Lamport 的逻辑时钟[3]非常相似,后者为生产者提供虚拟时钟。

在接收方处理无序消息: 这几乎是特定于应用程序的,但一般来说,当消息无序到达时,您有 2 个选择:a)丢弃较旧的消息,例如接收方必须显示 a 的最新值的情况库存。b) 有缓冲区来重新排序,就像在 TCP 连接中一样(例如,zookeeper 使用 TCP 作为 FIFO 排序的队列 [4-5])

工具: 如果您不为消息添加时间戳,则将所有消息从生产者按顺序发送到 Apache kafka分区,因为这将确保接收者可以按顺序接收消息。

如果您使用的消息系统不能保证有序传送(例如某些设置下的 AMQP [6]),那么您可以考虑为每条消息添加额外的单调递增数字/时钟。

[1] https://en.wiktionary.org/wiki/monotonic_increasing#targetText=形容词,contrast%20this%20with%20strictly%20increasing

[2] https://linux.die.net/man/2/clock_gettime

[3] https://en.wikipedia.org/wiki/Lamport_timestamps#Lamport 's_logic_clock_in_distributed_systems

[4] https://cwiki.apache.org/confluence/download/attachments/24193445/zookeeper-internals.pdf?version=1&modificationDate=1295034038000&api=v2

[5] http://www.tcs.hut.fi/Studies/T-79.5001/reports/2012-deSouzaMedeiros.pdf

[6] RabbitMQ - 消息传递顺序