【问题标题】:In C#, How can I get ALL the messages out of an Azure Service Bus Queue?在 C# 中,如何从 Azure 服务总线队列中获取所有消息?
【发布时间】: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


    【解决方案1】:

    一些基本假设:

    1. 请求消息的数量不保证是传递消息的数量。
    2. PeekLock 模式接收的消息将在某个时间点的锁定过期并被传递。

    如果您的目标是清除所有消息,您应该完成已收到的消息或在ReceiveAndDelete 模式下收到。这样你就不会再收到相同的消息了。如果您试图查看队列中的消息,那么您的LockDuration 需要足够长以确保所有消息都已被查看。

    我需要删除重复的消息并延迟处理一些“太新”的消息(由消息的最新副本确定),我想一次性处理服务总线队列的全部内容。

    更大的问题似乎是试图处理队列中的消息,就像数据库中的记录一样。重复检测已经是 Azure 服务总线的一项功能。延迟消息也是如此。但它需要一种不同于批处理的方法。

    【讨论】:

    • 我需要执行一些逻辑来确定消息是应该完成、放弃还是死信,所以我可能需要对 LockDuration 做一些事情,但我认为这可能会阻止我完成消息?我看到有一些关于重复检测的内容,但不知道如何为定义重复的内容提供逻辑。另外,我需要延迟哪些消息的条件逻辑,所以我不知道这是否可以帮助我......
    • 根据消息 ID 对消息进行重复数据删除。我为此写了一篇 [post](weblogs.asp.net/sfeldman/…)。也许它会帮助你。锁定持续时间不会阻止消息处理(完成、延迟、放弃)。锁定持续时间实际上是租用时间。长话短说,我强烈建议不要使用批处理方法。它只会让你受益。如果您需要消息会话,那就另当别论了。但不要强迫队列成为数据库。请改用 DB。
    • 这就是帖子显示的内容。阅读 l,然后回复 ? Cheers。
    • 这听起来像是您应该真正考虑审查的工作流程。为什么在一段时间内有重复的消息并且应该对最后一条消息采取措施?如果仅对最后一条消息采取行动,这些真的是“重复”吗?如果重复到达 5 分 10 秒会发生什么?处理重复是什么意思?不要认为 SO 会是一个很好的答案来源。如果这些消息确实是重复的并且可以安全地忽略,则处理第一条消息并根据最大时间范围对其余消息进行重复数据删除。
    • 如果消息不同,并且系列中的最后一条消息确实具有独特性(基于时间或消息数据),则需要采用不同的方法。一个可能需要一个状态。消息会话或带有“记忆元素”的东西,例如 saga(使用 NServiceBusMassTransit)都是一个选项。同样,我没有手头问题的完整背景。
    猜你喜欢
    • 2017-04-22
    • 1970-01-01
    • 2021-06-23
    • 2021-08-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多