【问题标题】:BlockingCollection<T>.TakeFromAny, for collections with a different generic typeBlockingCollection<T>.TakeFromAny,用于具有不同泛型类型的集合
【发布时间】:2016-11-19 01:07:41
【问题描述】:

.NET 中有一个BlockingCollection&lt;T&gt;.TakeFromAny 方法。它首先尝试快速获取 Take,然后默认为等待底层句柄的“慢”方法。我想用它来听上游生产者提供“消息”和下游生产者提供“结果”。

  • 是否可以使用 TakeFromAny - 或者是否有其他方法,无需重新实现此类 - 来监听异构类型阻塞集合集合的添加?

以下代码是类型有效的,自然无法编译:

object anyValue;
var collection = new List<BlockingCollection<object>>();
// following fails: cannot convert
//    from 'System.Collections.Concurrent.BlockingCollection<Message>'
//    to 'System.Collections.Concurrent.BlockingCollection<object>'
collection.Add(new BlockingCollection<Message>());
// fails for same reason
collection.Add(new BlockingCollection<Result>());
BlockingCollection<object>.TakeFromAny(collection.ToArray(), out anyValue);

可以只处理new BlockingCollection&lt;object&gt;() 实例并强制转换为 Take 以避免编译类型错误,尽管这让我犯了错误(呃) - 特别是因为通过方法接口丢失了输入。使用包装组合类型将解决后者; fsvo '解决'。


这里没有任何内容与问题直接相关,尽管它提供了上下文 - 对于那些感兴趣的人。提供核心基础设施功能的代码无法使用更高级别的构造(例如 Rx 或 TPL 数据流)。

这是一个基本的流程模型。生产者、代理和工作人员在不同的线程上运行(工作人员可以在同一个线程上运行,具体取决于任务调度程序的工作)。

[producer]   message -->   [proxy]   message --> [worker 1]
             <-- results             <-- results
                                     message --> [worker N..]
                                     <-- results

期望代理侦听消息(传入)和结果(返回)。代理执行一些工作,例如转换和分组,并将结果用作反馈。

将代理作为单独的线程将其与执行各种猴子业务的初始生产源隔离开来。工作任务用于并行性,而不是异步性,并且线程(在通过代理中的分组减少/消除争用之后)应该允许良好的扩展。

队列是在代理和工作人员之间建立的(而不是具有单个输入/结果的直接任务),因为在工作人员正在执行时,它可能会在结束之前处理额外的传入工作消息。这是为了确保工作人员可以延长/重用它在相关工作流中建立的上下文。

【问题讨论】:

  • 1.可以为每个集合阻塞一个线程吗? 2. 您是否需要有保证的“仅从一个集合中获取”操作,还是更像是“我想处理两个集合中的所有项目,但我不想并行执行”?
  • @svick 上游生产者(写入“消息”队列)和扇出下游消费者/生产者(他们在“结果”队列中产生结果)可以被阻止 - 这个“代理”消费者/生产者在与任何其他处理不同的线程上下文中运行,并将消息从上游移动到转换后的下游消费者,并将结果作为反馈移回上游。目标是“等待下一件要做的事情”,从任何队列中执行,然后再等待。
  • @svick 我想我最初的想法/设计是“投票/选择,使用类型化的阻塞集合”。
  • @svick 或者,“穷人版的演员”

标签: c# generics task-parallel-library blockingcollection


【解决方案1】:

我认为这里最好的选择是将两个阻塞集合的类型更改为您已经提到的BlockingCollection&lt;object&gt;,包括它的缺点。

如果您不能或不想这样做,另一种解决方案是合并 BlockingCollection&lt;object&gt; 并为每个源集合创建一个线程,将项目从其集合移动到合并的集合:

var producerCollection = new BlockingCollection<Message>();
var consumerCollection = new BlockingCollection<Results>();

var combinedCollection = new BlockingCollection<object>();

var producerCombiner = Task.Run(() =>
{
    foreach (var item in producerCollection.GetConsumingEnumerable())
    {
        combinedCollection.Add(item);
    }
});

var consumerCombiner = Task.Run(() =>
{
    foreach (var item in consumerCollection.GetConsumingEnumerable())
    {
        combinedCollection.Add(item);
    }
});

Task.WhenAll(producerCombiner, consumerCombiner)
    .ContinueWith(_ => combinedCollection.CompleteAdding());

foreach (var item in combinedCollection.GetConsumingEnumerable())
{
    // process item here
}

这不是很有效,因为它阻塞了两个额外的线程来执行此操作,但这是我能想到的最好的替代方案,不使用反射来访问TakeFromAny 使用的句柄。

【讨论】:

  • 我决定使用BlockingCollection&lt;object&gt; 路线,访问更加安全。在这个过程中,我发现 TakeFromAny 是相当有偏见的。它总是有利于第一个集合(如果他们有一个快速的接收,或者如果多个等待句柄同时发出信号,则在缓慢的接收)。
猜你喜欢
  • 2016-06-11
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-08-05
  • 1970-01-01
  • 2019-02-10
相关资源
最近更新 更多