【发布时间】:2021-07-26 03:47:27
【问题描述】:
因为我需要删除重复的消息并延迟处理一些“太新”的消息(由消息的最新副本确定),所以我想一次性处理服务总线队列的全部内容。
我不确定我可以期待多少条消息,但我非常乐观地认为它通常不应该是数百条,更不用说我认为可能是数以千计的限制了ReceiveAsync (int maxMessageCount, TimeSpan operationTimeout)。然而,事实证明,无论该值有多高,我在一次调用中只能读取大约 30 到 50 条消息
private async Task<IList<MicrosoftMessage>> Receive(IQueueConfig queueConfig) =>
await _messageReceiverLookup.GetMessageReceiver(queueConfig.QueueName)
.ReceiveAsync(queueConfig.MaximumRecords, TimeSpan.FromSeconds(10));
我尝试用一些额外的逻辑来包装它,比如:
List<MicrosoftMessage> messages = new();
List<MicrosoftMessage> newMessages = new();
do
{
newMessages = await ReceiveMessages(queueHandler, cancellationToken);
messages.AddRange(newMessages);
}
while (
newMessages.Count > 0
&& messages.Count > 0
&& messages.Count < queueHandler.QueueConfig.MaximumRecords
);
但发现这永远不会结束,因为系统会多次读取相同的消息。
然后我尝试了这个:
Dictionary<string, MicrosoftMessage> previosMessagesByToken;
Dictionary<string, MicrosoftMessage> allMessagesByToken = new();
List<MicrosoftMessage> newMessages;
do
{
previosMessagesByToken = allMessagesByToken;
newMessages = await ReceiveMessages(queueHandler, cancellationToken);
Dictionary<string, MicrosoftMessage> newMessagesByToken = newMessages.ToDictionary(x => x.SystemProperties.LockToken, x => x);
// Ensure we only collect each message once!
allMessagesByToken = allMessagesByToken.Concat(newMessagesByToken.Where(kvp => !allMessagesByToken.ContainsKey(kvp.Key)))
.ToDictionary(kvp => kvp.Key, kvp => kvp.Value);
}
while (
newMessages.Count > 0
&& allMessagesByToken.Count > previosMessagesByToken.Count
&& allMessagesByToken.Count < queueHandler.QueueConfig.MaximumRecords
);
这似乎可行,但一方面,我有一种直觉,这不应该那么复杂。另外,我并不完全相信这一点,因为我不完全理解为什么我没有收到所有消息,也没有收到重复的消息,所以我不禁觉得这个算法可能会允许一些消息落在裂缝,是不包括在内的非重复项。
有没有更好的方法可以让我获取所有消息?
【问题讨论】:
标签: c# message-queue bulk azure-servicebus-queues receiver