【发布时间】:2016-11-19 01:07:41
【问题描述】:
.NET 中有一个BlockingCollection<T>.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<object>() 实例并强制转换为 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