【问题标题】:How to yield from parallel tasks in .NET 4.5如何从 .NET 4.5 中的并行任务中获得收益
【发布时间】:2013-01-26 04:57:12
【问题描述】:

我想将 .NET 迭代器与并行任务/等待一起使用?。像这样的:

IEnumerable<TDst> Foo<TSrc, TDest>(IEnumerable<TSrc> source)
{
    Parallel.ForEach(
        source,
        s=>
        {
            // Ordering is NOT important
            // items can be yielded as soon as they are done                
            yield return ExecuteOrDownloadSomething(s);
        }
}

不幸的是 .NET 无法原生处理这个问题。 @svick 迄今为止的最佳答案 - 使用 AsParallel()。

奖励:任何实现多个发布者和单个订阅者的简单异步/等待代码?订阅者会屈服,而 pubs 会处理。 (仅限核心库)

【问题讨论】:

    标签: c# task-parallel-library async-await c#-5.0 yield-return


    【解决方案1】:

    这似乎是 PLINQ 的工作:

    return source.AsParallel().Select(s => ExecuteOrDownloadSomething(s));
    

    这将使用有限数量的线程并行执行委托,并在完成后立即返回每个结果。

    如果ExecuteOrDownloadSomething() 方法是IO 绑定的(例如它实际上下载了一些东西)并且你不想浪费线程,那么使用async-await 可能是有意义的,但它会更复杂。

    如果您想充分利用async,则不应返回IEnumerable,因为它是同步的(即,如果没有可用的项目,它会阻塞)。您需要的是某种异步收集,您可以使用 TPL Dataflow 中的 ISourceBlock(特别是 TransformBlock):

    ISourceBlock<TDst> Foo<TSrc, TDest>(IEnumerable<TSrc> source)
    {
        var block = new TransformBlock<TSrc, TDest>(
            async s => await ExecuteOrDownloadSomethingAsync(s),
            new ExecutionDataflowBlockOptions
            {
                MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded
            });
    
        foreach (var item in source)
            block.Post(item);
    
        block.Complete();
    
        return block;
    }
    

    如果源“慢”(即您想在迭代 source 完成之前开始处理来自 Foo() 的结果),您可能希望将 foreachComplete() 调用移动到单独的 @ 987654335@。更好的解决方案是将source 也变成ISourceBlock&lt;TSrc&gt;

    【讨论】:

    • 谢谢,但你能举例说明如何使用 async/await 解决这个问题吗?谢谢!
    • @Yurik 你能解释一下你为什么想要那个吗?
    • 主要是因为我觉得它可以帮助我理解新的 await 语法来解决不是“async 101”而是现实世界场景的问题。
    • 标记为已接受,但主要用于 AsParallel()。我想我真正想要的是使用 async/await 并仅使用核心库的多个 pubs + 单个 sub 的实现。不过谢谢!
    • @Yurik 为什么你关心一个库是否是核心?
    【解决方案2】:

    所以看来您真正想要做的是根据任务完成时间对一系列任务进行排序。这不是很复杂:

    public static IEnumerable<Task<T>> Order<T>(this IEnumerable<Task<T>> tasks)
    {
        var input = tasks.ToList();
    
        var output = input.Select(task => new TaskCompletionSource<T>());
        var collection = new BlockingCollection<TaskCompletionSource<T>>();
        foreach (var tcs in output)
            collection.Add(tcs);
    
        foreach (var task in input)
        {
            task.ContinueWith(t =>
            {
                var tcs = collection.Take();
                switch (task.Status)
                {
                    case TaskStatus.Canceled:
                        tcs.TrySetCanceled();
                        break;
                    case TaskStatus.Faulted:
                        tcs.TrySetException(task.Exception.InnerExceptions);
                        break;
                    case TaskStatus.RanToCompletion:
                        tcs.TrySetResult(task.Result);
                        break;
                }
            }
            , CancellationToken.None
            , TaskContinuationOptions.ExecuteSynchronously
            , TaskScheduler.Default);
        }
    
        return output.Select(tcs => tcs.Task);
    }
    

    所以在这里我们为每个输入任务创建一个TaskCompletionSource,然后遍历每个任务并设置一个延续,它从BlockingCollection 中获取下一个完成源并设置它的结果。完成的第一个任务获取返回的第一个 tcs,完成的第二个任务获取返回的第二个 tcs,依此类推。

    现在您的代码变得非常简单:

    var tasks = collection.Select(item => LongRunningOperationThatReturnsTask(item))
        .Order();
    foreach(var task in tasks)
    {
        var result = task.Result;//or you could `await` each result
        //....
    }
    

    【讨论】:

    • 谢谢,但我需要的是从方法中获取已处理对象的流作为产量。你提供的基本上是对 Parallel.ForEach() 的重写。
    • @Yurik 如果您不需要等待所有项目都完成,您可以删除WhenAll/WaitAll,但除此之外我看不到Select本身并不能满足您的需要。您有一系列项目,并且您希望将其转换为一系列任务,每个项目一个。 Select(item=&gt; LongRunningOperation(item)) 在返回一系列任务时如何不满足您的需求?
    • 在这种情况下,项目的顺序将与原始顺序相同,这可能是低效的。我不介意物品的无序产出。
    【解决方案3】:

    在 MS 机器人团队制作的异步库中,他们有并发原语,允许使用迭代器生成异步代码。

    图书馆 (CCR) 是免费的(它过去不是免费的)。一篇不错的介绍性文章可以在这里找到:Concurrent affairs

    也许您可以将这个库与 .Net 任务库一起使用,或者它会激发您“自己动手”的灵感

    【讨论】:

    • 您能解释一下您将如何在这里使用 CCR 吗?
    • 我引用的那篇文章比我能解释得更好。如果您查看它并检查该图:“图 6 SerialAsyncDemo”,它有一个代码示例,几乎与 OP 所要求的完全一样:使用 .Net 迭代器产生的异步操作。我承认我认为这种迭代器语法虽然在当时很聪明,但现在大多被 async/await 语法所取代
    猜你喜欢
    • 1970-01-01
    • 2012-09-02
    • 1970-01-01
    • 1970-01-01
    • 2020-02-26
    • 2015-06-21
    • 2013-12-19
    相关资源
    最近更新 更多