【问题标题】:How to correctly queue up tasks to run in C#如何正确排队任务以在 C# 中运行
【发布时间】:2015-11-25 23:19:33
【问题描述】:

我有一个项目枚举 (RunData.Demand),每个项目代表一些涉及通过 HTTP 调用 API 的工作。如果我只是 foreach 完成所有操作并在每次迭代期间调用 API,效果会很好。但是,每次迭代需要一两秒钟,所以我想运行 2-3 个线程并在它们之间分配工作。这就是我正在做的事情:

ThreadPool.SetMaxThreads(2, 5); // Trying to limit the amount of threads
var tasks = RunData.Demand
   .Select(service => Task.Run(async delegate
   {
      var availabilityResponse = await client.QueryAvailability(service);
      // Do some other stuff, not really important
   }));

await Task.WhenAll(tasks);

client.QueryAvailability 调用基本上是使用 HttpClient 类调用 API:

public async Task<QueryAvailabilityResponse> QueryAvailability(QueryAvailabilityMultidayRequest request)
{
   var response = await client.PostAsJsonAsync("api/queryavailabilitymultiday", request);

   if (response.IsSuccessStatusCode)
   {
      return await response.Content.ReadAsAsync<QueryAvailabilityResponse>();
   }

   throw new HttpException((int) response.StatusCode, response.ReasonPhrase);
}

这在一段时间内效果很好,但最终事情开始超时。如果我将 HttpClient Timeout 设置为一小时,我就会开始收到奇怪的内部服务器错误。

我开始做的是在QueryAvailability 方法中设置一个秒表来查看发生了什么。

正在发生的情况是 RunData.Demand 中的所有 1200 个项目同时被创建,并且所有 1200 个await client.PostAsJsonAsync 方法都被调用。然后它似乎使用 2 个线程来缓慢地检查任务,所以到最后我的任务已经等待了 9 或 10 分钟。

这是我想要的行为:

我想创建 1,200 个任务,然后在线程可用时一次运行 3-4 个。我确实不想立即排队 1,200 个 HTTP 调用。

有什么好的方法可以做到这一点吗?

【问题讨论】:

  • 您似乎没有为每次通话创建一个新的client。您知道 System.Net.Http.HttpClient 对于实例调用不是线程安全的吗?应该为每次调用创建(并在之后处理)一个新实例。
  • QueryAvailability 方法实际上位于创建HttpClient 的类中,HttpClient 是该实例的私有成员。虽然我不知道它不是线程安全的,但我绝对可以在每次调用之前创建它。我会进一步研究,谢谢!
  • 嗯,我做了一些研究,看来我正在做的是线程安全的。见here 和here

标签: c# .net multithreading asynchronous async-await


【解决方案1】:

正如我一直建议的那样。您需要的是 TPL 数据流(安装:Install-Package System.Threading.Tasks.Dataflow)。

您创建一个ActionBlock,其中包含要对每个项目执行的操作。设置MaxDegreeOfParallelism 进行节流。开始发帖并等待其完成:

var block = new ActionBlock<QueryAvailabilityMultidayRequest>(async service => 
{
    var availabilityResponse = await client.QueryAvailability(service);
    // ...
},
new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 4 });

foreach (var service in RunData.Demand)
{
    block.Post(service);
}

block.Complete();
await block.Completion;

【讨论】:

  • @MikeChristensen 是。它是 .Net 中为数不多的专门为 async-await 编写的库之一。
  • 完美运行!我现在是粉丝。接受。
【解决方案2】:

老问题,但我想提出一个使用 SemaphoreSlim 类的替代轻量级解决方案。只需引用 System.Threading。

SemaphoreSlim sem = new SemaphoreSlim(4,4);

foreach (var service in RunData.Demand)
{

    await sem.WaitAsync();
    Task t = Task.Run(async () => 
    {
        var availabilityResponse = await client.QueryAvailability(serviceCopy));    
        // do your other stuff here with the result of QueryAvailability
    }
    t.ContinueWith(sem.Release());
}

信号量充当锁定机制。您只能通过调用从计数中减去 1 的 Wait (WaitAsync) 来输入信号量。调用 release 将计数加一。

【讨论】:

  • 如果我理解正确,这取决于使用它的方式,可能适用于sem.Wait() 而不是await sem.WaitASync()。这样做会阻塞调用线程,因此不应在 UI 线程上完成,但在任何其他线程上,这可能是管理要完成的工作的最简单方法。具体来说,如果 4 个 sem 之一可用,它将立即进行。如果没有,它会等到有一个可用。
【解决方案3】:

您正在使用异步 HTTP 调用,因此限制线程数无济于事(Parallel.ForEach 中的ParallelOptions.MaxDegreeOfParallelism 也无济于事,正如答案之一所暗示的那样)。即使是单个线程也可以发起所有请求并在结果到达时对其进行处理。

解决它的一种方法是使用 TPL 数据流。

另一个不错的解决方案是将源IEnumerable 划分为多个分区,并按照this blog post 中的描述顺序处理每个分区中的项目:

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);
        }));
}

【讨论】:

    【解决方案4】:

    虽然 Dataflow 库很棒,但我认为不使用块组合时它有点重。我倾向于使用类似下面的扩展方法。

    此外,与 Partitioner 方法不同,它在调用上下文中运行异步方法 - 需要注意的是,如果您的代码不是真正异步的,或者采用“快速路径”,那么它将有效地同步运行,因为没有线程显式创建。

    public static async Task RunParallelAsync<T>(this IEnumerable<T> items, Func<T, Task> asyncAction, int maxParallel)
    {
        var tasks = new List<Task>();
    
        foreach (var item in items)
        {
            tasks.Add(asyncAction(item));
    
            if (tasks.Count < maxParallel)
                    continue; 
    
            var notCompleted = tasks.Where(t => !t.IsCompleted).ToList();
    
            if (notCompleted.Count >= maxParallel)
                await Task.WhenAny(notCompleted);
        }
    
        await Task.WhenAll(tasks);
    }
    

    【讨论】:

    • 在 foreach 循环的每次迭代中创建不必要的列表。与在每次迭代中计算未完成任务相比,用于限制的 ManualResetEvent 或 SempahoreSlim 将涉及更少的分配和更好的性能。
    • 感谢@Snak - 你是绝对正确的,实际上今天我使用了一种非常不同的辅助方法,没有任何特定的任务跟踪。但是请读者注意,使用简单的SemaphoreSlim 的幼稚方法很容易出现死锁 - 需要非常小心地确保任务异常正确传播以提前退出!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-02-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-06-16
    • 1970-01-01
    相关资源
    最近更新 更多