【问题标题】:Have a set of Tasks with only X running at a time有一组任务,一次只运行 X
【发布时间】:2012-12-28 19:58:03
【问题描述】:

假设我有 100 个任务需要 10 秒才能完成。 现在我想一次只运行 10 个,比如当这 10 个任务中的 1 个完成时,另一个任务会被执行,直到所有任务都完成。

现在我总是使用 ThreadPool.QueueUserWorkItem() 来完成此类任务,但我了解到这样做是不好的做法,我应该改用 Tasks。

我的问题是,我没有找到适合我的场景的好例子,所以您能否让我开始了解如何使用 Tasks 实现这一目标?

【问题讨论】:

  • 您从哪里了解到使用 ThreadPool 是不好的做法?
  • 我建议阅读一些文章和/或以前的 Stackoverflow 帖子,有很多其他人尝试过的编码示例以及提供答案的地方stackoverflow.com/questions/6192898/… 像我有 C# Stackoverflow ThreadPool.QueueUserWorkItem( )
  • 你想要一个阻塞直到所有任务完成的方法,还是你想要一个在所有任务完成后返回Task的方法?
  • 它应该阻塞,就像 ThreadPool 一样,但是有任务。 Stackoverflow 上的一些人在代码示例中告诉我 Threadpool 是不好的做法

标签: c# multithreading task


【解决方案1】:
SemaphoreSlim maxThread = new SemaphoreSlim(10);

for (int i = 0; i < 115; i++)
{
    maxThread.Wait();
    Task.Factory.StartNew(() =>
        {
            //Your Works
        }
        , TaskCreationOptions.LongRunning)
    .ContinueWith( (task) => maxThread.Release() );
}

【讨论】:

  • 为什么指定TaskCreationOptions.LongRunning
  • 如果您的任务是长期运行的任务,我建议您使用 Semaphore 类而不是 SemaphoreSlim。
  • @MarcGravell 我可能遗漏了一些东西,但我只看到一个主线程运行循环,并且工人(10)同时运行。被阻塞的是主线程(1)跨度>
  • @L.B.我道歉;我误读了等待的位置 - 我的错 - 我匆忙阅读,并且错误
  • @L.B 如何在 10 个任务中的任何一个任务完成时将新任务添加到队列中。表示将新任务添加到队列中。在这里,我的目标是,例如,如果我们已经通过 SemaphoreSlim 或 MaxDegreeOfParallelism 设置了一次运行 10 个任务的限制,但我不想创建 100 个任务,然后通过 SemaphoreSlim 或 MaxDegreeOfParallelism 设置限制并控制它们在 a单次。 ,我只想在10个任务中的任何一个任务完成后创建一个新任务,这个过程将无限继续。
【解决方案2】:

TPL Dataflow 非常适合做这样的事情。您可以轻松创建 Parallel.Invoke 的 100% 异步版本:

async Task ProcessTenAtOnce<T>(IEnumerable<T> items, Func<T, Task> func)
{
    ExecutionDataflowBlockOptions edfbo = new ExecutionDataflowBlockOptions
    {
         MaxDegreeOfParallelism = 10
    };

    ActionBlock<T> ab = new ActionBlock<T>(func, edfbo);

    foreach (T item in items)
    {
         await ab.SendAsync(item);
    }

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

【讨论】:

  • TPL 数据流库实际上非常酷,我已经找到了它的用途,谢谢指出
  • 这是一个很好的方法,但是否有可能让它返回? - 我正在考虑使用它来调用需要一段时间才能响应的服务。
  • @GabrielEspinoza 这只是 TPL Dataflow 可以做的一小部分。您可以使用TransformBlock 来满足您的需求。
  • TPL 数据流的 nuget 现已未列出。已替换为System.Threading.Tasks.Dataflow
  • 您能否更新此演示以显示用户应在何处添加命令以对每个项目执行,以及如何为项目传递参数?
【解决方案3】:

您有多种选择。初学者可以使用Parallel.Invoke

public void DoWork(IEnumerable<Action> actions)
{
    Parallel.Invoke(new ParallelOptions() { MaxDegreeOfParallelism = 10 }
        , actions.ToArray());
}

这是一个替代选项,它会更加努力地运行恰好 10 个任务(尽管线程池中处理这些任务的线程数可能不同)并返回一个 Task 指示它何时完成,而不是阻塞直到完成。

public Task DoWork(IList<Action> actions)
{
    List<Task> tasks = new List<Task>();
    int numWorkers = 10;
    int batchSize = (int)Math.Ceiling(actions.Count / (double)numWorkers);
    foreach (var batch in actions.Batch(actions.Count / 10))
    {
        tasks.Add(Task.Factory.StartNew(() =>
        {
            foreach (var action in batch)
            {
                action();
            }
        }));
    }

    return Task.WhenAll(tasks);
}

如果你没有 MoreLinq,对于 Batch 函数,这是我更简单的实现:

public static IEnumerable<IEnumerable<T>> Batch<T>(this IEnumerable<T> source, int batchSize)
{
    List<T> buffer = new List<T>(batchSize);

    foreach (T item in source)
    {
        buffer.Add(item);

        if (buffer.Count >= batchSize)
        {
            yield return buffer;
            buffer = new List<T>();
        }
    }
    if (buffer.Count >= 0)
    {
        yield return buffer;
    }
}

【讨论】:

  • Now I want to only run 10 at a time, MaxDegreeOfParallelism 只是一个上限。
  • @L.B 好吧,从技术上讲,如果工作单元的数量不能被 10 整除,那么 正好 10 是不可能的。
  • 不,即使有 100 个作品,它最终可能会运行更少的任务。它取决于许多参数,例如 # of CPUs
  • @L.B 我知道这一点。我是说,即使你试图接近 10,你也不能总是完美的,你只能......更接近。
  • @L.B 我添加的附加版本你满意吗?
【解决方案4】:

我很想使用我能想到的最简单的解决方案,就像我认为使用 TPL 一样:

string[] urls={};
Parallel.ForEach(urls, new ParallelOptions() { MaxDegreeOfParallelism = 2}, url =>
{
   //Download the content or do whatever you want with each URL
});

【讨论】:

    【解决方案5】:

    你可以像这样创建一个方法:

    public static async Task RunLimitedNumberAtATime<T>(int numberOfTasksConcurrent, 
        IEnumerable<T> inputList, Func<T, Task> asyncFunc)
    {
        Queue<T> inputQueue = new Queue<T>(inputList);
        List<Task> runningTasks = new List<Task>(numberOfTasksConcurrent);
        for (int i = 0; i < numberOfTasksConcurrent && inputQueue.Count > 0; i++)
            runningTasks.Add(asyncFunc(inputQueue.Dequeue()));
    
        while (inputQueue.Count > 0)
        {
            Task task = await Task.WhenAny(runningTasks);
            runningTasks.Remove(task);
            runningTasks.Add(asyncFunc(inputQueue.Dequeue()));
        }
    
        await Task.WhenAll(runningTasks);
    }
    

    然后你可以调用任何异步方法 n 次,限制如下:

    Task task = RunLimitedNumberAtATime(10,
        Enumerable.Range(1, 100),
        async x =>
        {
            Console.WriteLine($"Starting task {x}");
            await Task.Delay(100);
            Console.WriteLine($"Finishing task {x}");
        });
    

    或者如果你想运行长时间运行的非异步方法,你可以这样做:

    Task task = RunLimitedNumberAtATime(10,
        Enumerable.Range(1, 100),
        x => Task.Factory.StartNew(() => {
            Console.WriteLine($"Starting task {x}");
            System.Threading.Thread.Sleep(100);
            Console.WriteLine($"Finishing task {x}");
        }, TaskCreationOptions.LongRunning));
    

    也许框架的某处有类似的方法,但我还没找到。

    【讨论】:

    • 我已经修改了这段代码来运行不同的无参数任务,效果很好
    猜你喜欢
    • 2022-01-09
    • 1970-01-01
    • 2013-01-12
    • 1970-01-01
    • 1970-01-01
    • 2011-09-21
    • 2018-09-05
    • 1970-01-01
    相关资源
    最近更新 更多