【问题标题】:awaiting async method in consumer using BlockingCollection as queue在使用 BlockingCollection 作为队列的消费者中等待异步方法
【发布时间】:2016-07-15 16:15:53
【问题描述】:

我正在开发一个服务器端控制台应用程序,该应用程序从多个 WCF 服务接收数据,对其进行一些处理,然后使用 SignalR 通过单个连接将结果转发到 IIS 服务器。

我尝试使用生产者消费者模式来实现这一点,其中 WCF 服务是生产者,使用 SignalR 发送数据的类是消费者。对于队列,我使用了BlockingCollection。

但是,当使用 await/async 在消费者的 while 循环中发送数据时,会卡住,直到所有其他线程完成将数据添加到队列中。

出于测试目的,我已经用Task.Delay(1000).Wait(); 或await Task.Delay(1000); 替换了实际发送数据的代码,它们也都卡住了。 一个简单的Thread.Sleep(1000); 似乎工作得很好,这让我认为异步代码是问题所在。

所以我的问题是:有什么东西阻止异步代码在 while 循环中完成吗?我错过了什么?

我正在像这样启动消费者线程:

new Thread(Worker).Start();

以及消费者代码:

private void Worker()
{
    while (!_queue.IsCompleted)
    {
        IMobileMessage msg = null;
        try
        {
            msg = _queue.Take();
        }
        catch (InvalidOperationException)
        {
        }

        if (msg != null)
        {
            try
            {
                Trace.TraceInformation("Sending: {0}", msg.Name);
                Thread.Sleep(1000); // <-- works
                //Task.Delay(1000).Wait(); // <-- doesn't work
                msg.SentTime = DateTime.UtcNow;
                Trace.TraceInformation("X sent at {1}: {0}", msg.Name, msg.SentTime);
            }
            catch (Exception e)
            {
                TraceException(e);
            }
        }
    }
}

【问题讨论】:

  • 阻止和异步不是朋友。如果您将async 与BlockingCollection&lt;T&gt; 混合使用,您应该强烈考虑放弃BlockingCollection&lt;T&gt; 并查看TPL 数据流。 BufferBlock&lt;T&gt; 是一个很好的起点,大致相当于BlockingCollection&lt;T&gt;,但数据流可以为生产者/消费者场景提供更多功能。花时间去了解它。这是值得的。
  • 好的,非常感谢,我会调查的。

标签: c# .net multithreading asynchronous blockingcollection


【解决方案1】:

正如 spender 正确指出的那样,BlockingCollection(顾名思义)仅适用于阻塞代码,不适用于异步代码。

有异步兼容的生产者/消费者队列,例如BufferBlock&lt;T&gt;。在这种情况下,我认为ActionBlock&lt;T&gt; 会更好:

private ActionBlock<IMobileMsg> _block = new ActionBlock<IMobileMsg>(async msg =>
{
  try
  {
    Trace.TraceInformation("Sending: {0}", msg.Name);
    await Task.Delay(1000);
    msg.SentTime = DateTime.UtcNow;
    Trace.TraceInformation("X sent at {1}: {0}", msg.Name, msg.SentTime);
  }
  catch (Exception e)
  {
    TraceException(e);
  }
});

这将替换您的整个消费线程和主循环。

【讨论】:

  • 我添加了您建议替换我的消费者和队列的代码,并阅读了您在博客上关于 TPL 数据流的优秀文章,它看起来有效,但我现在无法确认我的生产者计时器每次 2 个刻度后停止。但我想这是新的东西。一旦我确认它有效,我会尽快接受你的回答!谢谢!
  • 由于某种原因,我似乎无法让它工作......我的生产者每秒检查是否有来自 WCF 服务的新消息(在循环中使用 Task.Delay(1000); 并将它们发布到ActionBlock(如果有)。根据我在操作块的延迟和生产者的延迟中使用的时间,一切正常,完全停止,或者只将消息发布到块缓冲区但不再处理它们。任何明显的原因你不能在生产者中使用带有异步 while 循环的操作块?
  • @J.Neijt:不,异步生产者工作正常。如果您的 TraceException 抛出异常,则该块将停止工作。
  • 不幸的是我无法解决这个问题,随着时间的流逝,我不得不回到以前的版本,我现在已经改变了架构。可能还有其他一些问题正在发生......所以我无法确认,但我可以看到你的答案是所描述问题的解决方案,我很乐意接受它并非常感谢你的帮助!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-05-09
  • 2013-04-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-11-17
  • 1970-01-01
相关资源
最近更新 更多