【问题标题】:Limiting the amount of concurrent tasks in .NET 4.5在 .NET 4.5 中限制并发任务的数量
【发布时间】:2013-12-19 19:00:56
【问题描述】:

观察以下函数:

public Task RunInOrderAsync<TTaskSeed>(IEnumerable<TTaskSeed> taskSeedGenerator,
    CreateTaskDelegate<TTaskSeed> createTask,
    OnTaskErrorDelegate<TTaskSeed> onError = null,
    OnTaskSuccessDelegate<TTaskSeed> onSuccess = null) where TTaskSeed : class
{
    Action<Exception, TTaskSeed> onFailed = (exc, taskSeed) =>
    {
        if (onError != null)
        {
            onError(exc, taskSeed);
        }
    };

    Action<Task> onDone = t =>
    {
        var taskSeed = (TTaskSeed)t.AsyncState;
        if (t.Exception != null)
        {
            onFailed(t.Exception, taskSeed);
        }
        else if (onSuccess != null)
        {
            onSuccess(t, taskSeed);
        }
    };

    var enumerator = taskSeedGenerator.GetEnumerator();
    Task task = null;
    while (enumerator.MoveNext())
    {
        if (task == null)
        {
            try
            {
                task = createTask(enumerator.Current);
                Debug.Assert(ReferenceEquals(task.AsyncState, enumerator.Current));
            }
            catch (Exception exc)
            {
                onFailed(exc, enumerator.Current);
            }
        }
        else
        {
            task = task.ContinueWith((t, taskSeed) =>
            {
                onDone(t);
                var res = createTask((TTaskSeed)taskSeed);
                Debug.Assert(ReferenceEquals(res.AsyncState, taskSeed));
                return res;
            }, enumerator.Current).TaskUnwrap();
        }
    }

    if (task != null)
    {
        task = task.ContinueWith(onDone);
    }

    return task;
}

其中TaskUnwrap 是标准Task.Unwrap 的状态保留版本:

public static class Extensions
{
    public static Task TaskUnwrap(this Task<Task> task, object state = null)
    {
        return task.Unwrap().ContinueWith((t, _) =>
        {
            if (t.Exception != null)
            {
                throw t.Exception;
            }
        }, state ?? task.AsyncState);
    }
}

RunInOrderAsync 方法允许异步运行 N 个任务,但顺序是 - 一个接一个。实际上,它运行从给定种子创建的任务,并发限制为 1。

让我们假设createTask委托从种子创建的任务不对应于多个并发任务。

现在,我想输入 maxConcurrencyLevel 参数,因此函数签名如下所示:

Task RunInOrderAsync<TTaskSeed>(int maxConcurrencyLevel,
  IEnumerable<TTaskSeed> taskSeedGenerator,
  CreateTaskDelegate<TTaskSeed> createTask,
  OnTaskErrorDelegate<TTaskSeed> onError = null,
  OnTaskSuccessDelegate<TTaskSeed> onSuccess = null) where TTaskSeed : class

在这里我有点卡住了。

SO 有如下问题:

基本上提出了两种解决问题的方法:

  1. Parallel.ForEachParallelOptions 一起使用,指定MaxDegreeOfParallelism 属性值等于所需的最大并发级别。
  2. 使用具有所需 MaximumConcurrencyLevel 值的自定义 TaskScheduler

第二种方法并没有减少它,因为涉及的所有任务都必须使用相同的任务调度程序实例。为此,所有用于返回 Task 的方法都必须具有接受自定义 TaskScheduler 实例的重载。不幸的是,微软在这方面并不是很一致。例如,SqlConnection.OpenAsync 不接受这样的论点(但 TaskFactory.FromAsync 接受)。

第一种方法意味着我必须将任务转换为操作,如下所示:

() => t.Wait()

我不确定这是一个好主意,但我很乐意就此获得更多意见。

另一种方法是使用TaskFactory.ContinueWhenAny,但这很麻烦。

有什么想法吗?

编辑 1

我想澄清想要限制的原因。我们的任务最终对同一个 SQL 服务器执行 SQL 语句。我们想要的是一种限制并发传出 SQL 语句数量的方法。完全有可能在其他代码段中同时执行其他 SQL 语句,但这是一个批处理器,可能会淹没服务器。

现在,请注意,虽然我们讨论的是同一个 SQL 服务器,但同一台服务器上有许多数据库。所以,这并不是要限制打开同一个数据库的 SQL 连接的数量,因为数据库可能根本不一样。

这就是为什么像 ThreadPool.SetMaxThreads() 这样的末日解决方案是无关紧要的。

现在,关于SqlConnection.OpenAsync。它之所以异步是有原因的——它可能会往返于服务器,因此可能会受到网络延迟和分布式环境的其他可爱副作用的影响。因此,它与接受TaskScheduler 参数的其他异步方法没有什么不同。我倾向于认为不接受只是一个错误。

编辑 2

我想保留原始函数的异步精神。因此,我希望避免任何明确的阻塞解决方案。

编辑 3

感谢@fsimonazzi's answer 我现在有了所需功能的有效实现。代码如下:

        var sem = new SemaphoreSlim(maxConcurrencyLevel);
        var tasks = new List<Task>();

        var enumerator = taskSeedGenerator.GetEnumerator();
        while (enumerator.MoveNext())
        {
            tasks.Add(sem.WaitAsync().ContinueWith((_, taskSeed) =>
            {
                Task task = null;
                try
                {
                    task = createTask((TTaskSeed)taskSeed);
                    if (task != null)
                    {
                        Debug.Assert(ReferenceEquals(task.AsyncState, taskSeed));
                        task = task.ContinueWith(t =>
                        {
                            sem.Release();
                            onDone(t);
                        });
                    }
                }
                catch (Exception exc)
                {
                    sem.Release();
                    onFailed(exc, (TTaskSeed)taskSeed);
                }
                return task;
            }, enumerator.Current).TaskUnwrap());
        }

        return Task.Factory.ContinueWhenAll(tasks.ToArray(), _ => sem.Dispose());

【问题讨论】:

  • BlockingCollection 允许你设置一个 BoundedCapacity
  • 好吧,当然 SqlConnection.OpenAsync 没有那个选项。它不会烧掉 N 个线程。事实上,它不消耗任何的可能性很大。如果您的任务非常不守规矩,以至于您必须提供这种保证,那么 ThreadPool.SetMaxThreads() 是核选项。
  • 请参阅EDIT 1
  • @Blam - 如果您提供代码,您的回复将是一个很好的答案。
  • 这就是为什么它是评论而不是答案。请参阅 MSDN 上的文档。有一个使用 BoundedCapacity 的示例。 SO 不是代码生成工具。

标签: .net asynchronous


【解决方案1】:

您可以使用信号量来限制处理。使用 WaitAsync() 方法,您可以获得预期的异步。像这样的东西(为简洁起见,删除了错误处理):

private static async Task DoStuff<T>(int maxConcurrency, IEnumerable<T> items, Func<T, Task> createTask)
{
    using (var sem = new SemaphoreSlim(maxConcurrency))
    {
        var tasks = new List<Task>();

        foreach (var item in items)
        {
            await sem.WaitAsync();
            var task = createTask(item).ContinueWith(t => sem.Release());
            tasks.Add(task);
        }

        await Task.WhenAll(tasks);
    }
}

已编辑以消除在所有发布操作有机会执行之前可以释放信号量的错误。

【讨论】:

  • 最后的WhenAll可以换成ContinueWhenAll吗?
  • 如果你有一个延续,他们可以使用 ContinueWhenAll。
  • 谢谢。结果证明您的方法很容易实现 - 在 EDIT 3 中发布了代码。
【解决方案2】:

这里已经有很多答案了。我想解决您在斯蒂芬斯回答中所做的评论,关于使用 TPL 数据流限制并发的示例。即使你在这个问题的另一个答案中留下了评论说你不再使用基于任务的方法,它可能会帮助其他人。

为此使用ActionBlock&lt;T&gt; 的示例是:

private static async Task DoStuff<T>(int maxConcurrency, IEnumerable<T> items, Func<T, Task> createTask)
{
    var ab = new ActionBlock<T>(createTask, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = maxConcurrency });

    foreach (var item in items)
    {
        ab.Post(item);
    }

    ab.Complete();
    await ab.Completion;
}

关于 TPL 数据流的更多信息可以在这里找到:https://msdn.microsoft.com/en-us/library/system.threading.tasks.dataflow(v=vs.110).aspx

【讨论】:

    【解决方案3】:

    目前可用的两个最佳解决方案是Semaphoreslim(根据@fsimonazzi's answer)和TPL 数据流块(即ActionBlock&lt;T&gt;TransformBlock&lt;T&gt;)。这两个区块都有simple way to set the level of concurrency

    Parallel 不是一种理想的方法,因为您需要阻塞异步操作,为每个操作使用一个线程池线程。

    另外,TaskScheduler 在这里也不起作用。仅供参考,TaskScheduler 通过async 方法“继承”的,正如我在my async intro blog post 上描述的那样。它对您的问题不起作用的原因是因为任务调度程序只控制 executing 任务,而不是 event 任务 - 所以,像OpenAsync 这样的 SQL 操作不会“ count”计入并发限制。

    【讨论】:

    • 您介意使用 TPL 数据流块绘制方法吗?
    • 查看我使用ActionBlock&lt;T&gt;的答案。
    【解决方案4】:

    这是@fsimonazzi 答案的变体,没有 SemaphoreSlim,虽然很酷。

    private static async Task DoStuff<T>(int maxConcurrency, IEnumerable<T> items, Func<T, Task> createTask)
    {
        var tasks = new List<Task>();
        foreach (var item in items)
        {
            if (tasks.Count >= maxConcurrency)
            {
                await Task.WhenAll(tasks);
                tasks.Clear();
            }
            var task = createTask(item);
            tasks.Add(task);
        }
        await Task.WhenAll(tasks);
    }
    

    【讨论】:

      【解决方案5】:

      这是@scott-turner 答案的变体,虽然很酷。他的回答是以 maxConcurrency 块提交工作,并等待每个块完全完成,然后再提交下一个块。此变体根据需要提交新任务,以尝试并确保 maxConcurrency 任务始终处于运行状态。它还演示了如何使用 Task 而不是 Task。

      请注意,与 SemaphoreSlim 变体相比,它的好处是使用 SemaphoreSlim 您需要等待两种不同类型的任务 - 信号量和工作。如果工作的类型是 Task 而不是 Task,那就有问题了。

          private static async Task<R[]> concurrentAsync<T, R>(int maxConcurrency, IEnumerable<T> items, Func<T, Task<R>> createTask)
          {
              var allTasks = new List<Task<R>>();
              var activeTasks = new List<Task<R>>();
              foreach (var item in items)
              {
                  if (activeTasks.Count >= maxConcurrency)
                  {
                      var completedTask = await Task.WhenAny(activeTasks);
                      activeTasks.Remove(completedTask);
                  }
                  var task = createTask(item);
                  allTasks.Add(task);
                  activeTasks.Add(task);
              }
              return await Task.WhenAll(allTasks);
          }
      

      【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-02-23
      • 1970-01-01
      • 1970-01-01
      • 2022-10-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多