【问题标题】:Can this kind of concurrency problem be solved with async/await?这种并发问题可以用 async/await 解决吗?
【发布时间】:2020-02-07 21:13:15
【问题描述】:

我有这样的功能:

static void AddResultsToDb(IEnumerable<int> numbers)
{
    foreach (int number in numbers)
    {
        int result = ComputeResult(number); // This takes a long time, but is thread safe.
        AddResultToDb(number, result); // This is quick but not thread safe.
    }
}

我可以解决这个问题,例如,使用Parallel.ForEach 计算结果,然后使用常规foreach 将结果添加到数据库中。

但是,出于教育目的,我想要一个围绕等待/异步的解决方案。但是,无论我读了多少关于它的内容,我都无法将其全神贯注。如果 await/async 在这种情况下不适用,我想了解原因。

【问题讨论】:

  • 既然目的是教育性的,我们是否可以假设调用的两个方法是异步的(ComputeResultAsyncAddResultToDbAsync)?因为如果它们不是异步的,那么这个练习的教育价值将非常小,如果不是负面的话。 async/await 是一种旨在解决异步问题而不是并发问题的技术。

标签: c# concurrency async-await


【解决方案1】:

正如其他人所建议的那样,这不是使用async/await 的情况,因为那是异步的。你正在做的是并发。微软有一个专门为此而设计的框架,它很好地解决了这个问题。

因此,出于学习目的,您应该使用 Microsoft 的响应式框架(又名 Rx)- NuGet System.Reactive 并添加 using System.Reactive.Linq; - 然后您可以这样做:

static void AddResultsToDb(IEnumerable<int> numbers)
{
    numbers 
        .ToObservable()
        .SelectMany(n => Observable.Start(() => new { n, r = ComputeResult(n) }))
        .Do(x => AddResultToDb(x.n, x.r))
        .Wait();
}

SelectMany/Observable.Start 组合允许尽可能多的ComputeResult 调用同时发生。 Rx 的好处是它会序列化结果,这样一次只有一个调用会转到AddResultToDb


要控制并行度,您可以将 SelectMany 更改为 Select/Merge,如下所示:

static void AddResultsToDb(IEnumerable<int> numbers)
{
    numbers 
        .ToObservable()
        .Select(n => Observable.Start(() => new { n, r = ComputeResult(n) }))
        .Merge(maxConcurrent: 2)
        .Do(x => AddResultToDb(x.n, x.r))
        .Wait();
}

【讨论】:

  • RX 已进入聊天室...不错
  • 这种方法的问题,抛开房间里的大象(相当庞大而复杂的RX.Net库的巨大学习曲线)不谈,它既不是阻塞也不是异步的。您无法等到结果添加到数据库中。在这方面,它类似于一劳永逸的Task。解决这个问题并非易事。
  • @TheodorZoulias - 这很简单。将.Subscribe(x =&gt; AddResultToDb(x.n, x.r)) 更改为.Do(x =&gt; AddResultToDb(x.n, x.r)).Wait()
  • 不错!我撤销了我的反对票。是否可以控制计算的并行度,并按初始顺序保存结果?
  • @TheodorZoulias - 我添加了更改并行度的代码。要保存原始顺序(OP 没有要求您注意),您可以在开头使用 .Select((x, n) =&gt; 来记录项目的顺序,然后在 Do 之前使用对 ToArray() 的调用然后在使用SelectMany 展开之前对所有元素进行排序。
【解决方案2】:

async 和 await 模式并不真正适合您的第一种方法。它非常适合IO Bound 工作负载以实现可扩展性,或者适合具有 UI 响应能力的框架。它不太适合原始 CPU 工作负载

但是您仍然可以从并行处理中获益,因为您的第一种方法昂贵并且线程安全

在下面的示例中,我使用Parallel LINQ (PLINQ) 来流畅地表达结果,而无需担心预先确定大小的数组 / 并发集合 / 锁定,尽管您可以使用其他 TPL 功能,例如 Parallel.For/ForEach

// Potentially break up the workloads in parallel
// return the number and result in a ValueTuple
var results = numbers.AsParallel()
                     .Select(x => (number: x, result: ComputeResult(x)))
                     .ToList();

// iterate through the number and results and execute them serially 
foreach (var (number, result) in results)
   AddResultToDb(number, result);

注意:这里的假设是顺序不重要


补充

您的方法AddResultToDb 看起来只是将结果插入到数据库 中,这是IO Bound 并且值得async,此外可能会获取所有结果一次并将它们插入 bulk/batch 节省 往返行程


来自评论信用@TheodorZoulias

preserve the order,你可以使用AsOrdered的方法,在 一些性能损失的代价。可能的表现 改进是去掉ToList(),这样结果就加起来了 与计算同时发送到数据库。

为了尽快提供结果,这可能是一个不错的选择 禁用happens by default 的部分缓冲的想法,通过 链接方法 查询中的.WithMergeOptions(ParallelMergeOptions.NotBuffered)

var results = numbers.AsParallel()
                     .Select(x => (number: x, result: ComputeResult(x)))
                     .WithMergeOptions(ParallelMergeOptions.NotBuffered)
                     .AsOrdered();

示例


其他资源

ParallelEnumerable.AsOrdered Method

启用对数据源的处理,就好像它是有序的一样,覆盖 无序的默认值。 AsOrdered 只能在非泛型上调用 序列

ParallelEnumerable.WithMergeOptions

设置此查询的合并选项,指定查询的方式 将缓冲输出。

ParallelMergeOptions Enum

NotBuffered 使用没有输出缓冲区的合并。计算出结果元素后,立即将该元素提供给 查询的消费者。

【讨论】:

  • preserve the order,您可以使用AsOrdered 方法,但会降低性能。一个可能的性能改进是删除ToList(),以便在计算的同时将结果添加到数据库中。
  • @TheodorZoulias hah 别再跟踪我了!哦,图片是一样的
  • 你做得太简单了?!与删除ToList 相关的更多细节。为了使结果尽快可用,最好通过在查询中链接方法 .WithMergeOptions(ParallelMergeOptions.NotBuffered) 来禁用部分缓冲 that happens by default
【解决方案3】:

async/await 的情况并非如此,因为听起来ComputeResult 的计算成本很高,而不是只需要很长的、不确定的时间。 aync/await 更适合您真正等待的任务。 Parallel.ForEach 实际上会处理您的工作负载。

如果有的话,AddResultToDb 是您想要异步/等待的 - 您将等待外部操作完成。

好深入的解释:https://stackoverflow.com/a/35485780/127257

【讨论】:

    【解决方案4】:

    老实说,使用Parallel.For 似乎是最简单的解决方案,因为您的计算可能会受 CPU 限制。 Async/await 更适合 I/O 绑定操作,因为它不需要另一个线程等待 I/O 操作完成(请参阅there is no thread)。

    话虽如此,您仍然可以将 async/await 用于放置在线程池中的任务。所以你可以这样做。

    static void AddResultToDb(int number)
    {
        int result = ComputeResult(number); 
        AddResultToDb(number, result); 
    }
    
    static async Task AddResultsToDb(IEnumerable<int> numbers)
    {
        var tasks = numbers.Select
        (
            number => Task.Run( () => AddResultToDb(number) )  
        )
        .ToList();
    
        await Task.WhenAll(tasks);
    }
    

    【讨论】:

    • 应该是await AddResultToDb(number, result)我猜。但如果AddResultToDb 不是线程安全的,这是否明智?还是我错过了早间咖啡
    • @MichaelRandall 你说得对,要么我需要添加一个await,要么从AddResultToDb() 中删除async Task。我认为这些操作是线程安全的,因为 OP 正在询问 Parallel.For 的替代方案。
    • @JohnWu - OP 明确表示 AddResultToDb 不是线程安全的。
    猜你喜欢
    • 2016-01-09
    • 2020-03-25
    • 1970-01-01
    • 2012-02-07
    • 2011-07-16
    • 1970-01-01
    • 2011-10-11
    • 2023-03-19
    相关资源
    最近更新 更多