【问题标题】:Returning a BlockingCollection as IEnumerable from a method从方法返回 BlockingCollection 作为 IEnumerable
【发布时间】:2012-02-22 04:51:03
【问题描述】:

我正在尝试从 BlockingCollection 支持的方法返回 IEnumerable。代码模式是:

public IEnumerable<T> Execute() {   
    var results = new BlockingCollection<T>(10);  
    _ExecuteLoad(results);   
    return results.GetConsumingEnumerable(); 
}

private void _ExecuteLoad<T>(BlockingCollection<T> results) {
    var loadTask = Task.Factory.StartNew(() =>
    { 
        //some async code that adds items to results
        results.CompleteAdding();
    });
}

public void Consumer() {
    var count = Execute().Count();
}

问题是从 Execute() 返回的可枚举总是空的。我见过的所有示例都在任务中迭代 BlockingCollection。在这种情况下,这似乎行不通。

有人知道我哪里出错了吗?


为了让事情更清楚一点,我粘贴了我正在执行以填充集合的代码。也许这里有什么导致问题的原因?

Task.Factory.StartNew(() =>
{
    var continuationRowKey = "";
    var continuationParitionKey = "";
    var action = HttpMethod.Get;
    var queryUri = _GetTableQueryUri(tableServiceUri, tableName, query, continuationParitionKey, continuationRowKey, timeout);
    while (true)
    {
        using (var request = GetRequest(queryUri, null, action.Method, azureAccountName, azureAccountKey))
        {
            request.Method = action;
            request.RequestUri = queryUri;

            using (var client = new HttpClient())
            {
                var sendTask = client.SendAsync(request, HttpCompletionOption.ResponseHeadersRead);
                using (var response = sendTask.Result)
                {
                    continuationParitionKey = // stuff from headers
                    continuationRowKey = // stuff from headers

                    var streamTask = response.Content.ReadAsStreamAsync();
                    using (var stream = streamTask.Result)
                    {
                        using (var reader = XmlReader.Create(stream))
                        {
                            while (reader.Read())
                            {
                                if (reader.NodeType == XmlNodeType.Element && reader.Name == "entry" && reader.NamespaceURI == "http://www.w3.org/2005/Atom")
                                {
                                    results.Add(XNode.ReadFrom(reader) as XElement);
                                }
                            }
                            reader.Close();
                        }
                    }
                }
            }

            if (continuationParitionKey == null && continuationRowKey == null)
                break;

            queryUri = _GetTableQueryUri(tableServiceUri, tableName, query, continuationParitionKey, continuationRowKey, timeout);
        }
    }
    results.CompleteAdding();
});

【问题讨论】:

  • 它并不能解决问题,但是如果您只是要立即使用Result,那么调用...Async() 方法没有多大意义。
  • 您的代码原始代码对我来说很好用。您是否尝试过调试Task 中的代码,看看为什么它从不调用Add()?因为这是最有可能的解释。

标签: .net multithreading asynchronous task-parallel-library async-ctp


【解决方案1】:

完成向集合添加项目后,您需要调用results.CompleteAdding()

如果不这样做,枚举将永远不会结束,Count() 将永远不会返回。

除此之外,您发布的代码是正确的。

【讨论】:

  • 我应该补充一下。问题是方法马上返回,枚举总是空的。
  • Count() 将阻塞直到枚举完成。您是否尝试在CompleteAdding() 上设置断点?
  • 感谢您的回答。事实证明,我有一个不相关的错误导致 CompleteAdding 没有被调用。谢谢!
猜你喜欢
  • 1970-01-01
  • 2023-03-31
  • 1970-01-01
  • 1970-01-01
  • 2016-12-10
  • 2010-09-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多