【问题标题】:Using Parallel foreach in an API request process?在 API 请求过程中使用并行 foreach?
【发布时间】:2019-11-29 12:06:58
【问题描述】:

在处理请求时使用Parallel.ForEach 是否正确?我问这个是因为 async task 旨在尽可能多地参加会议,而不是尽可能快,这是 Parallel.ForEach 会做的。

简单示例:

public async Task<OperationResult> ProcessApiRequest(List<string> ids)
{
    Parallel.ForEach(ids, async (id) =>
            {         
                await this.doStuff(id);
                await this.doAnotherStuff(id);
            });

    return OperationResult.Success();
}

想象一下,我可以收到 1 个或 100 万个 id,并且我希望尽可能多地参与请求。由于我的线程将忙于处理 100 万个 id,它将难以处理与会者的新请求,对吗?

谢谢!

【问题讨论】:

  • Parallel.ForEach 用于数据并行,而不是并发操作。此呼叫将触发 1M async void 呼叫,没有人会等待。它几乎会立即返回,也许在任何请求有机会开始之前。编译器还会发出警告,ProcessApiRequest 没有等待,将同时运行。
  • 是的,这正是我所担心的。谢谢!

标签: c# api asynchronous .net-core parallel.foreach


【解决方案1】:

不,这是不正确的。 Parallel.ForEach 用于数据并行。它将创建与机器上的核心一样多的工作任务,对输入数据进行分区并每个分区使用一个工作人员。它对async 操作一无所知,这意味着您的代码本质上是:

Parallel.ForEach(ids, async void (int id) =>
        {         
            await this.doStuff(id);
            await this.doAnotherStuff(id);
        });

在四台机器上,这将触发 100 万个请求,一次 4 个,无需等待任何一个。它可以在任何请求有机会完成之前轻松返回。

如果您想以受控方式执行多个请求,您可以使用具有特定并行度的 ActionBlock,例如:

var options=new ExecutionDataflowBlockOptions
{
    MaxDegreeOfParallelism = 10,
    BoundedCapacity=100
}
var block=new ActionBlock<string>(async id=>{....},options);


foreach(var id in ids)
{
    await block.SendAsync(id);
}
block.Complete();
await block.Completion;

该块将处理多达 10 个并发请求。如果动作真的是异步的,或者异步等待时间很长,我们可以轻松地使用比可用内核数更高的 DOP。

输入消息是缓冲的,这意味着我们最终可能会在慢块的输入缓冲区中等待 1M 请求。为避免这种情况,BoundedCapacity 设置将在 SendAsync 无法接受更多输入时阻止。

最后,对Complete() 的调用告诉块我们已经完成,它应该处理其输入缓冲区中的所有剩余消息。我们等待他们完成await block.Completion

【讨论】:

  • 就我而言,调用者不需要等待await this.doStuff(id) 或其他任何操作,我返回的是Acepted 状态码。所以在这里我根本不需要Parallel.ForEach,因为我只需要处理“东西”。谢谢你的好解释。
【解决方案2】:

您的担心是对的,Parallel.ForEach 默认情况下会使用线程池中尽可能多的线程,线程池将逐渐扩大到它需要的最大线程数。 Task.Run 对于 Web 服务器来说通常是个坏主意,Parallel.ForEach 通常会差很多倍。

特别是考虑到ids 是无限的,您可能很快就会遇到这样一种情况,即您的请求将被排队,因为所有线程都忙于满足少数请求。

所以你的担心是对的,这种代码正在优化单个请求的延迟以实现非常低的规模,但在规模上会牺牲一个公平且性能良好的网络服务器,最终消除延迟初始延迟的胜利,并创建你是一个更广泛的服务问题。

更新 - 正如 Panagiotis Kanavos 在 cmets 中指出的那样,Parallel.ForEach 没有 Task 重载,因此只会运行委托的初始同步部分,而留下大部分异步工作排队,您的 API 刚刚着火,可能在不知不觉中忘记了。

对于使用 ChannelReaderChannelWriter 的完全异步生产者消费者模式的替代版本,以及一些新的 C# 8.0 语法,您可以试试这个:

public async Task<OperationResult> ProcessApiRequest(List<string> ids)
{
    var channel = Channel.CreateBounded<string>(new BoundedChannelOptions(100) {SingleWriter = true});

    foreach (var id in ids)
    {
        await channel.Writer.WriteAsync(id); // If the back pressure exceeds 100 ids, we asynchronously wait here
    }
    channel.Writer.Complete();

    for (var i = 0; i < 8; i++) // 8 concurrent readers
    {
        _ = Task.Run(async () =>
        {
            await foreach (var id in channel.Reader.ReadAllAsync())
            {
                await this.doStuff(id);
                await this.doAnotherStuff(id);
            }
        });
    }

    return OperationResult.Success();
}

【讨论】:

  • 这不正确。 Parallel.ForEach 将创建与机器上的内核一样多的工作任务。它适用于数据并行,而不是并发操作,因此使用更多是没有意义的。 HUGE 问题是 Parallel.ForEach 对 async 一无所知,因此 所有 这些调用本质上都是 async void 调用。他们将尽可能快地被解雇,一次 4 或 8 个,而无需等待其中任何一个完成
  • 我没有发现缺少异步委托重载,我会更新我的答案。你能指出我支持声称它是核心数量的文档吗,其他 api 这样做,但是这个使用默认的 TaskScheduler 文档说这将是无限的。
  • 它与 API 或 TaskScheduler 无关。没有接受 Task&lt;&gt;ForEach 重载
  • 它说的是Behind the scenes, the Task Scheduler partitions the task based on system resources and workload. When possible, the scheduler redistributes work among multiple threads and processors if the workload becomes unbalanced. 它根本没有提到使用 20 亿个线程。它说您可以自定义分区和调度,但如果您不这样做,明智的默认设置是使用所有内核。
  • 不,它将使用机器中的所有 核心,而不是该线程池中所有 1000 个可能的线程。除非您显式增加了 DOP,或者使用了强制调度程序增加工作人员数量的阻塞代码。所有这些都在文档中进行了描述
猜你喜欢
  • 1970-01-01
  • 2021-07-16
  • 1970-01-01
  • 2019-08-15
  • 2019-02-24
  • 2011-10-15
  • 2013-01-16
  • 2017-05-03
  • 2022-01-06
相关资源
最近更新 更多