【问题标题】:Using SemaphoreSlim with Parallel.ForEach将 SemaphoreSlim 与 Parallel.ForEach 一起使用
【发布时间】:2020-11-30 15:38:44
【问题描述】:

这就是我想要达到的目标。假设我有一个每分钟运行一次并执行一些 I/O 操作的进程。我希望 5 个线程同时执行并执行操作。假设如果 2 个线程花费的时间超过一分钟,当进程在一分钟后再次运行时,它应该同时执行 3 个线程,因为 2 个线程已经在执行一些操作。

所以,我使用了SemaphoreSlimParallel.ForEach 的组合。请让我知道这是实现此目的的正确方法还是有其他更好的方法。

private static SemaphoreSlim _semaphoreSlim = new SemaphoreSlim(5);

private async Task ExecuteAsync()
{
    try
    {
        var availableThreads = _semaphoreSlim.CurrentCount;

        if (availableThreads > 0)
        {
            var lists = await _feedSourceService.GetListAsync(availableThreads); // select @top(availableThreads) * from table

            Parallel.ForEach(
                lists,
                new ParallelOptions
                {
                    MaxDegreeOfParallelism = availableThreads
                },
                async item =>
                {
                    await _semaphoreSlim.WaitAsync();

                    try
                    {
                        // I/O operations
                    }
                    finally
                    {
                        _semaphoreSlim.Release();
                    }
                });
        }
    }
    catch (Exception ex)
    {
        _logger.LogError(ex.Message, ex);
    }
}

【问题讨论】:

  • 所有使用asyncParallel 的代码都不正确。
  • @StephenCleary 你宁愿使用await Task.WhenAll(tasks);吗?
  • @StephenCleary 你能否解释一下为什么它是错误的。这将有助于理解
  • @ReyanChougle 我会看here
  • @ReyanChougle 你仍然可以使用 Semaphor 进行节流

标签: asp.net multithreading task-parallel-library semaphore parallel.foreach


【解决方案1】:

假设我有一个每分钟运行一次并执行一些 I/O 操作的进程...假设如果 2 个线程花费的时间超过一分钟,并且当该进程在一分钟后再次运行时,它应该同时执行 3 个线程作为 2线程已经在做一些操作了。

这种问题描述有些常见,但很难正确编码。这是因为您有一个轮询式计时器(基于时间),它试图定期调整一个节流机制。正确地做到这一点非常困难。

所以,我建议的第一件事是更改问题描述。考虑让轮询机制读取所有未完成的工作,然后从那里使用正常的限制(例如,将 then 添加到执行受限的ActionBlock)。

也就是说,如果您希望继续走更复杂的路径,这样的代码可以避免 Parallelasync 的问题:

private static SemaphoreSlim _semaphoreSlim = new SemaphoreSlim(5);

private async Task ExecuteAsync()
{
    try
    {
        var availableThreads = _semaphoreSlim.CurrentCount;

        if (availableThreads > 0)
        {
            var lists = await _feedSourceService.GetListAsync(availableThreads); // select @top(availableThreads) * from table

            var tasks = lists.Select(
                async item =>
                {
                    await _semaphoreSlim.WaitAsync();

                    try
                    {
                        // I/O operations
                    }
                    finally
                    {
                        _semaphoreSlim.Release();
                    }
                }).ToList();
            await Task.WhenAll(tasks);
        }
    }
    catch (Exception ex)
    {
        _logger.LogError(ex.Message, ex);
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多