【问题标题】:BrokeredMessage Automatically Disposed after calling OnMessage()BrokeredMessage 调用 OnMessage() 后自动释放
【发布时间】:2015-05-26 20:05:27
【问题描述】:

我正在尝试对 Azure 服务总线中的项目进行排队,以便批量处理它们。我知道 Azure 服务总线有一个 ReceiveBatch() 但它似乎有问题,原因如下:

  • 我一次最多只能收到 256 条消息,即使这样也可以根据消息大小随机发送。
  • 即使我偷看有多少消息正在等待,我也不知道要进行多少 RequestBatch 调用,因为我不知道每次调用会返回多少消息。由于消息会不断涌入,因此我不能在它为空之前继续发出请求,因为它永远不会为空。

我决定只使用消息监听器,它比浪费窥视更便宜,并且会给我更多的控制权。

基本上我试图让一定数量的消息建立起来 然后立即处理它们。我使用计时器来强制延迟,但我需要 以便能够在我的物品进入时对其进行排队。

根据我的计时器要求,阻塞集合似乎不是一个好的选择,所以我正在尝试使用 ConcurrentBag。

var batchingQueue = new ConcurrentBag<BrokeredMessage>();
myQueueClient.OnMessage((m) =>
{
    Console.WriteLine("Queueing message");
    batchingQueue.Add(m);
});

while (true)
{
    var sw = WaitableStopwatch.StartNew();
    BrokeredMessage msg;
    while (batchingQueue.TryTake(out msg)) // <== Object is already disposed
    {
        ...do this until I have a thousand ready to be written to DB in batch
        Console.WriteLine("Completing message");
        msg.Complete(); // <== ERRORS HERE
    }

    sw.Wait(MINIMUM_DELAY);
}

但是,一旦我在 OnMessage 之外访问消息 管道它显示 BrokeredMessage 已被释放。

我认为这一定是 OnMessage 的某种自动行为,除了立即处理我不想做的事情之外,我看不到任何方法可以对消息做任何事情。

【问题讨论】:

    标签: c# multithreading azure queue servicebus


    【解决方案1】:

    BlockingCollection 非常容易做到这一点。

    var batchingQueue = new BlockingCollection<BrokeredMessage>();
    
    myQueueClient.OnMessage((m) =>
    {
        Console.WriteLine("Queueing message");
        batchingQueue.Add(m);
    });
    

    还有你的消费者线程:

    foreach (var msg in batchingQueue.GetConsumingEnumerable())
    {
        Console.WriteLine("Completing message");
        msg.Complete();
    }
    

    GetConsumingEnumerable 返回一个迭代器,它消耗队列中的项目,直到设置了IsCompleted 属性并且队列为空。如果队列为空但IsCompletedFalse,则它会进行非忙等待下一项。

    要取消消费者线程(即关闭程序),您停止向队列添加内容并让主线程调用batchingQueue.CompleteAdding。消费者将队列清空,看到IsCompleted属性为True,然后退出。

    在这里使用BlockingCollectionConcurrentBagConcurrentQueue 更好,因为BlockingCollection 接口更易于使用。特别是,GetConsumingEnumerable 的使用使您不必担心检查计数或忙等待(轮询循环)。它只是工作。

    还要注意ConcurrentBag 有一些相当奇怪的删除行为。特别是,删除项目的顺序会根据删除项目的线程而有所不同。创建袋子的线程以与其他线程不同的顺序移除项目。详情请见Using the ConcurrentBag Collection

    您还没有说明为什么要在输入时对项目进行批处理。除非有压倒一切的性能原因这样做,否则使用批处理逻辑使代码复杂化似乎不是一个特别好的主意。


    如果您想批量写入数据库,那么我建议使用简单的List&lt;T&gt; 来缓冲项目。如果您必须在将项目写入数据库之前对其进行处理,请使用我上面展示的技术来处理它们。然后,与其直接写入数据库,不如将项目添加到列表中。当列表获得 1,000 个项目或经过给定的时间量时,分配一个新列表并启动一个任务以将旧列表写入数据库。像这样:

    // at class scope
    
    // Flush every 5 minutes.
    private readonly TimeSpan FlushDelay = TimeSpan.FromMinutes(5);
    private const int MaxBufferItems = 1000;
    
    // Create a timer for the buffer flush.
    System.Threading.Timer _flushTimer = new System.Threading.Timer(TimedFlush, FlushDelay.TotalMilliseconds, Timeout.Infinite);  
    
    // A lock for the list. Unless you're getting hundreds of thousands
    // of items per second, this will not be a performance problem.
    object _listLock = new Object();
    
    List<BrokeredMessage> _recordBuffer = new List<BrokeredMessage>();
    

    然后,在您的消费者中:

    foreach (var msg in batchingQueue.GetConsumingEnumerable())
    {
        // process the message
        Console.WriteLine("Completing message");
        msg.Complete();
        lock (_listLock)
        {
            _recordBuffer.Add(msg);
            if (_recordBuffer.Count >= MaxBufferItems)
            {
                // Stop the timer
                _flushTimer.Change(Timeout.Infinite, Timeout.Infinite);
    
                // Save the old list and allocate a new one
                var myList = _recordBuffer;
                _recordBuffer = new List<BrokeredMessage>();
    
                // Start a task to write to the database
                Task.Factory.StartNew(() => FlushBuffer(myList));
    
                // Restart the timer
                _flushTimer.Change(FlushDelay.TotalMilliseconds, Timeout.Infinite);
            }
        }
    }
    
    private void TimedFlush()
    {
        bool lockTaken = false;
        List<BrokeredMessage> myList = null;
    
        try
        {
            if (Monitor.TryEnter(_listLock, 0, out lockTaken))
            {
                // Save the old list and allocate a new one
                myList = _recordBuffer;
                _recordBuffer = new List<BrokeredMessage>();
            }
        }
        finally
        {
            if (lockTaken)
            {
                Monitor.Exit(_listLock);
            }
        }
    
        if (myList != null)
        {
            FlushBuffer(myList);
        }
    
        // Restart the timer
        _flushTimer.Change(FlushDelay.TotalMilliseconds, Timeout.Infinite);
    }
    

    这里的想法是,您将旧列表移开,分配一个新列表以便继续处理,然后将旧列表的项目写入数据库。锁是为了防止计时器和记录计数器相互踩踏。如果没有锁,事情可能会在一段时间内正常运行,然后您会在不可预知的时间发生奇怪的崩溃。

    我喜欢这种设计,因为它消除了消费者的轮询。我唯一不喜欢的是消费者必须知道计时器(即它必须停止然后重新启动计时器)。稍加思考,我就可以消除这个要求。但它的编写方式效果很好。

    【讨论】:

    • 批处理用于对数据库进行批量写入,这比每个请求一进来就写入它们要快得多。我之前看过 BlockingCollection 但我不确定如何实现我的计时器逻辑用它。我需要让队列建立一秒钟,除非我有 1000 个等待去,然后我继续写入数据库。我也不确定它是否解决了我提到的垃圾收集问题。我会看看你的代码并玩一下。
    • 这段代码有同样的问题: BrokeredMessage 已被释放。
    • 我不喜欢从单独的线程中删除 BrokeredMessage。如果我以相同的方法将其删除,则将其添加到队列中,它可以正常工作。在我尝试将其从处理循环中删除的那一刻,代理消息显示为已处理。
    • @KingOfHypocrites:我怀疑您的消息正在被其他人处理。您确定它没有被生产者中的某些东西处理吗? BlockingCollectionConcurrentQueue 都不会对集合中的项目进行任何更改。我建议在对象的 Dispose 方法上放置一个断点,这样你就可以看到它在哪里被杀死了。
    • 是的,我很肯定。一旦我尝试将消息从它处理的队列中拉出。我示例中的代码几乎就是整个代码示例。如果我在 OnMessage 线程上将项目从队列中拉出,然后调用 Complete() 它工作正常。只要我在代码示例中显示的 while 循环中执行相同的操作,BrokerMessage 就会被释放。
    【解决方案2】:

    切换到 OnMessageAsync 为我解决了问题

    _queueClient.OnMessageAsync(async receivedMessage =>
    

    【讨论】:

      【解决方案3】:

      我就 MSDN 上的 BrokeredMessage 被处理问题联系了 Microsoft,这是回复:

      非常基本的规则,我不确定这是否记录在案。接收到的消息需要在回调函数的生命周期内进行处理。在您的情况下,消息将在异步回调完成时被释放,这就是您的完整尝试失败并在另一个线程中出现 ObjectDisposedException 的原因。

      我真的不明白排队消息以进行进一步处理对吞吐量有何帮助。这肯定会给客户增加更多的负担。尝试在异步回调中处理消息,这应该足够了。

      在我的情况下,这意味着我不能以我想要的方式使用 ServiceBus,我必须重新考虑我希望事情如何工作。混蛋。

      【讨论】:

        【解决方案4】:

        我在开始使用 Azure 服务总线服务时遇到了同样的问题。

        我发现 OnMessage 方法总是处理 BrokedMessage 对象。 Jim Mischel 提出的方法对我没有帮助(但读起来很有趣 - 谢谢!)。

        经过一番调查,我发现整个方法是错误的。让我解释做你想做的事情的正确方法。

        1. 仅在 OnMessage 方法处理程序中使用 BrokedMessage.Complete() 方法。
        2. 如果您需要在此方法之外处理消息,您应该使用方法 QueueClient.Complete(Guid lockToken)。 “LockToken”是 BrokeredMessage 对象的属性。

        例子:

         var messageOptions = new OnMessageOptions {
              AutoComplete       = false,
              AutoRenewTimeout   = TimeSpan.FromMinutes( 5 ),
             MaxConcurrentCalls = 1
         };
         var buffer = new Dictionary<string, Guid>();
        
         // get message from queue 
         myQueueClient.OnMessage(
              m => buffer.Add(key: m.GetBody<string>(), value: m.LockToken), 
              messageOptions // this option says to ServiceBus to "froze" message in he queue until we process it
         );         
        
         foreach(var item in buffer){
            try {
                Console.WriteLine($"Process item: {item.Key}");
                myQueueClient.Complete(item.Value);// you can also use method CompleteBatch(...) to improve performance
            } 
            catch{
                // "unfroze" message in ServiceBus. Message would be delivered to other listener 
                myQueueClient.Defer(item.Value);
            }
         }
        

        【讨论】:

          【解决方案5】:

          我的解决方案是获取消息 SequenceNumber,然后推迟消息并将 SequenceNumber 添加到 BlockingCollection。一旦 BlockingCollection 选择了一个新项目,它就可以通过 SequenceNumber 接收延迟消息并将消息标记为完成。如果由于某种原因 BlockingCollection 不处理 SequenceNumber,它将保留在队列中作为延迟,以便稍后在重新启动进程时将其拾取。如果在 BlockingCollection 中仍有项目时进程异常终止,这可以防止丢失消息。

          BlockingCollection<long> queueSequenceNumbers = new BlockingCollection<long>();
          
          //This finds any deferred/unfinished messages on startup. 
          BrokeredMessage existingMessage = client.Peek();
          while (existingMessage != null)
          {
              if (existingMessage.State == MessageState.Deferred)
              {
                  queueSequenceNumbers.Add(existingMessage.SequenceNumber);
              }
              existingMessage = client.Peek();
          }
          
          
          //setup the message handler
          Action<BrokeredMessage> processMessage = new Action<BrokeredMessage>((message) =>
          {
              try
              {
                  //skip deferred messages if they are already in the queueSequenceNumbers collection.
                  if (message.State != MessageState.Deferred || (message.State == MessageState.Deferred && !queueSequenceNumbers.Any(x => x == message.SequenceNumber)))
                  {
                      message.Defer();
                      queueSequenceNumbers.Add(message.SequenceNumber);
                  }
          
              }
              catch (Exception ex)
              {
                   // Indicates a problem, unlock message in queue
                   message.Abandon();
              }
          
          });
          
          
          // Callback to handle newly received messages
          client.OnMessage(processMessage, new OnMessageOptions() { AutoComplete = false, MaxConcurrentCalls = 1 });            
          
          //start the blocking loop to process messages as they are added to the collection
          foreach (var queueSequenceNumber in queueSequenceNumbers.GetConsumingEnumerable())
          {
               var message = client.Receive(queueSequenceNumber);
               //mark the message as complete so it's removed from the queue
               message.Complete();                 
               //do something with the message       
          }
          

          【讨论】:

            猜你喜欢
            • 1970-01-01
            • 2014-02-03
            • 1970-01-01
            • 1970-01-01
            • 2011-11-09
            • 2011-05-03
            • 1970-01-01
            • 2013-11-19
            相关资源
            最近更新 更多