【问题标题】:Throttling with SemaphoreSlim -- "Task.Run()" vs "new Func<Task>()"使用 SemaphoreSlim 进行节流——“Task.Run()”与“new Func<Task>()”
【发布时间】:2018-07-09 15:09:28
【问题描述】:

这可能不是专门针对 SemaphoreSlim,但基本上我的问题是关于以下两种限制长时间运行任务集合的方法之间是否存在差异,如果是,那么差异是什么(以及何时有使用任何一个)。

在下面的示例中,假设每个跟踪的任务都涉及从 Url 加载数据(完全是虚构的示例,但这是我在 SemaphoreSlim 示例中找到的常见示例)。

主要区别在于如何将单个任务添加到跟踪任务列表中。在第一个示例中,我们使用 lambda 调用 Task.Run(),而在第二个示例中,我们使用 lambda 新建一个 Func(&lt;Task&lt;Result&gt;&gt;()),然后立即调用该 func 并将结果添加到跟踪的任务列表中。

示例:

使用Task.Run():

 SemaphoreSlim ss = new SemaphoreSlim(_concurrentTasks);
 List<string> urls = ImportUrlsFromSource();

 List<Task<Result>> trackedTasks = new List<Task<Result>>();
        foreach (var item in urls)
        {
            await ss.WaitAsync().ConfigureAwait(false);
            trackedTasks.Add(Task.Run(async () =>
            {

                try
                {
                    return await ProcessUrl(item);
                }
                catch (Exception e)
                {
                    _log.Error($"logging some stuff");
                    throw;
                }
                finally
                {
                    ss.Release();
                }
            }));
        }
        var results = await Task.WhenAll(trackedTasks);

使用新的函数:

 SemaphoreSlim ss = new SemaphoreSlim(_concurrentTasks);
 List<string> urls = ImportUrlsFromSource();

 List<Task<Result>> trackedTasks = new List<Task<Result>>();
        foreach (var item in urls)
        {
            trackedTasks.Add(new Func<Task<Result>>(async () =>
            {
                await ss.WaitAsync().ConfigureAwait(false);
                try
                {
                    return await ProcessUrl(item);
                }
                catch (Exception e)
                {
                    _log.Error($"logging some stuff");
                    throw;
                }
                finally
                {
                    ss.Release();
                }
            })());
        }
        var results = await Task.WhenAll(trackedTasks);

【问题讨论】:

  • 不重新发明轮子而只使用 PLinQ:var results = urls.AsParallel().WithDegreeOfParallelism(_concurrentTasks).Select(SyncDownLoadMethod).ToList() 可能会容易得多。
  • 很公平。但我仍然对这里的基本问题感到好奇——这些示例的工作方式的实际区别是什么?

标签: c# asynchronous async-await task


【解决方案1】:

有两个区别:


Task.Run 进行错误处理

首先,当您调用 lambda 时,它会运行。另一方面,Task.Run 会调用它。这是相关的,因为Task.Run 在幕后做了一些工作。它的主要工作是处理一个错误的任务......

如果您调用 lambda,而 lambda 抛出,它会在您将 Task 添加到列表之前抛出...

但是,在您的情况下,由于您的 lambda 是异步的,编译器将为它创建 Task(您不是手动创建的),它会正确处理异常并通过返回的 @ 使其可用987654326@。因此这一点没有实际意义


Task.Run 防止任务附件

Task.Run 设置DenyChildAttach。这意味着在Task.Run 中创建的任务独立于返回的Task 运行(不与之同步)。

例如这段代码:

List<Task<int>> trackedTasks = new List<Task<int>>();
var numbers = new int[]{0, 1, 2, 3, 4};
foreach (var item in numbers)
{
    trackedTasks.Add(Task.Run(async () =>
    {
        var x = 0;
        (new Func<Task<int>>(async () =>{x = item; return x;}))().Wait();
        Console.WriteLine(x);
        return x;
    }));
}
var results = await Task.WhenAll(trackedTasks);

将以未知的顺序输出 0 到 4 的数字。但是下面的代码:

List<Task<int>> trackedTasks = new List<Task<int>>();
var numbers = new int[]{0, 1, 2, 3, 4};
foreach (var item in numbers)
{
    trackedTasks.Add(new Func<Task<int>>(async () =>
    {
        var x = 0;
        (new Func<Task<int>>(async () =>{x = item; return x;}))().Wait();
        Console.WriteLine(x);
        return x;
    })());
}
var results = await Task.WhenAll(trackedTasks);

每次都会按顺序输出从 0 到 4 的数字。这很奇怪,对吧?发生的事情是内部任务附加到外部任务,并立即在同一个线程中执行。但是如果你使用Task.Run,内部任务不会被独立附加和调度。

即使您使用await,这仍然适用,只要您await 的任务不转到外部系统...

外部系统会发生什么?好吧,例如,如果您的任务正在从 URL 读取 - 就像在您的示例中一样 - 系统将创建一个 TaskCompletionSource,从中获取 Task,设置一个将结果写入 TaskCompletionSource 的响应处理程序,提出请求,并返回Task。这个Task 没有计划,它与父任务在同一个线程上运行是没有意义的。因此,它可以打破秩序。

由于您使用await 等待外部系统,这一点也没有实际意义


结论

我必须得出结论,它们是等价的。

如果您想确保安全,并确保它按预期工作,即使 - 在未来的版本中 - 上述某些要点不再有实际意义,请保留 Task.Run。另一方面,如果您真的想优化,请使用 lambda 并避免 Task.Run(非常小的)开销。不过,这可能不会成为瓶颈。


附录

当我谈论到外部系统的任务时,我指的是在 .NET 之外运行的东西。有一些代码将在 .NET 中运行以与外部系统交互,但大部分代码不会在 .NET 中运行,因此根本不会在托管线程中。

API 的使用者没有指定任何事情发生。该任务将是一个承诺任务,但它没有公开,对于消费者来说并没有什么特别之处。

事实上,到外部系统的任务可能几乎不会在 CPU 中运行。此外,它可能只是在等待计算机外部的某些东西(可能是网络或用户输入)。

模式如下:

  1. 库创建TaskCompletionSource

  2. 库设置了接收通知的方法。它可以是回调、事件、消息循环、钩子、监听套接字、管道、等待全局互斥体......任何需要的东西。

  3. 库设置代码以响应通知,该通知将调用SetResultSetException 上的TaskCompletionSource 作为收到的通知的适当位置。

  4. 库对外部系统进行实际调用。

  5. 库返回TaskCompletionSource.Task

注意:特别注意优化,不要在不应该​​重新排序的地方,并注意在设置阶段处理错误。此外,如果涉及CancellationToken,则必须将其考虑在内(并在适当时在TaskCompletionSource 上调用SetCancelled)。此外,对通知的反应(或取消)可能需要拆除。 啊,别忘了验证你的参数。

然后,外部系统开始执行它所做的任何事情。然后当它完成或出现问题时,向库发出通知,并且您的 Task 突然完成,出现故障......(或者如果发生取消,您的 Task 现在被取消)并且.NET 将安排继续根据需要完成任务。

注意:async/await 在幕后使用延续,这就是执行恢复的方式。

顺便说一句,如果你想自己实现 SempahoreSlim,你必须做一些与我上面描述的非常相似的事情。你可以在我的backport of SemaphoreSlim看到它。


让我们看几个 promise 任务的例子......

  1. Task.Delay:当我们用Task.Delay 等待时,CPU 没有旋转。这不是在线程中运行。在这种情况下,通知机制将是一个操作系统计时器。当 OS 看到计时器的时间已经过去时,它会调用 CLR,然后 CLR 会将任务标记为已完成。什么线程在等待?没有。

  2. FileStream.ReadSync:当我们使用FileStream.ReadSync从存储中读取时,实际工作由设备完成。 CRL 必须声明一个自定义事件,然后将事件、文件句柄和缓冲区传递给操作系统……操作系统调用设备驱动程序,设备驱动程序与设备交互。当存储设备恢复信息时,它会通过 DMA 技术写入内存(直接在指定的缓冲区上)。完成后,它将设置一个中断,由驱动程序处理,通知操作系统,调用自定义事件,将任务标记为已完成。哪个线程从存储中读取数据?没有。

将使用类似的模式从网页下载,但这次设备连接到网络。如何发出 HTTP 请求以及系统如何等待响应超出了本答案的范围。

也有可能外部系统是另一个程序,在这种情况下它会在线程上运行。但它不会是您进程中的托管线程。


您的收获是这些任务不会在您的任何线程上运行。他们的时间可能取决于外部因素。因此,将它们视为在同一个线程中运行或我们可以预测它们的时间是没有意义的(当然,除了定时器的情况)。

【讨论】:

  • 这是一个非常有用的解释。快速跟进以确保我理解 - 使用外部系统和异步调用,返回的任务未安排这一事实仅意味着它可以在任何线程上运行并且不同步,对吗?而这个事实(在外部系统中的非计划/有序运行)是“自动的”,因为它不需要我指定它就发生了,对吗?
  • 我想基本上我想知道您是否可以扩展您的意思“此任务未计划,它与父任务在同一线程上运行没有意义。因此,它可以破坏秩序。”
  • @istrupin 扩展答案。希望它能解决您的疑虑。
【解决方案2】:

两者都不是很好,因为它们会立即创建任务。 func 版本的开销要少一些,因为它将Task.Run 路由保存在线程池上,只是为了立即结束线程池工作并在信号量上挂起。您不需要 async Func,您可以通过使用 async 方法(可能是本地函数)来简化它。

但你根本不应该这样做。相反,请使用helper method that implements a parallel async foreach

public static Task ForEachAsync<T>(this IEnumerable<T> source, int dop, Func<T, Task> body) 
{ 
    return Task.WhenAll( 
        from partition in Partitioner.Create(source).GetPartitions(dop) 
        select Task.Run(async delegate { 
            using (partition) 
                while (partition.MoveNext()) 
                    await body(partition.Current); 
        })); 
}

那你就去urls.ForEachAsync(myDop, async input =&gt; await ProcessAsync(input));

在这里,任务是按需创建的。您甚至可以使输入流变得惰性。

【讨论】:

  • 只是为了确认一下——我使用并行异步 foreach 而不是使用带有信号量的异步方法在性能方面有什么优势?是不是我不再花费资源来实际创建任务?我想确保我明白真正的好处就是全部。
  • 是的,信号量版本创建任务。信号量还维护一个任务队列,以便接下来激活。此外,所有这些机器只是不必要的代码,而且很复杂。很难理解和验证是否正确。你想要一个 foreach,使用一个 foreach。不幸的是,框架没有它,所以您需要将这个小助手复制到每个新项目中。
猜你喜欢
  • 1970-01-01
  • 2022-11-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多