You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.NET 6异步Semaphore高负载下报错及代码使用问题咨询

问题分析:SemaphoreSlim高负载下抛出超量错误的原因与修复

我正在开发一个基础的非数据库连接池,要求每个项目仅允许创建1个连接。该连接池需支持异步任务/多线程环境,因此使用Semaphore而非常规Lock。

低负载时测试代码可正常运行,但高负载下Semaphore会抛出错误:

Adding the specified count to the semaphore would cause it to exceed its maximum count

降低负载(比如增大_waitTimeMs或减少_numberOfTasks)可解决该错误,偶尔高负载测试也能正常执行。请问代码或Semaphore的使用是否存在错误或误解?

using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;

internal class Program
{
    static int _numberOfTasks = 50;
    static int _waitTimeMs = 1;
    static SemaphoreSlim _dictLock = new SemaphoreSlim(1, 1);
    static ConcurrentDictionary<string, bool> _pool = new ConcurrentDictionary<string, bool>();

    /// <summary>
    /// 每个项目仅允许1个连接。
    /// 若连接池中有可用连接则复用,否则创建1个新连接。
    /// </summary>
    static async Task<string> GetConnection(string projId)
    {
        try
        {
            // 进入信号量锁,防止同一项目添加多个连接
            if (await _dictLock.WaitAsync(_waitTimeMs))
            {
                // 尝试从连接池获取连接
                if (_pool.TryGetValue(projId, out bool value))
                {
                    if (value == false)
                        return "Exists but not connected yet.";
                    else
                        return "Success, exists and connected.";
                }
                // 否则将连接添加到池
                else
                {
                    _pool.TryAdd(projId, false);

                    // 模拟建立新连接的延迟
                    await Task.Delay(2);

                    _pool.TryUpdate(projId, true, false);

                    return "Created new connection successfully & added to pool.";
                }
            }
            // 报告未能及时获取锁的失败信息
            else
                return "Server busy. Please try again later.";
        }
        catch (Exception ex)
        {
            return "Error " + ex.Message;
        }
        finally
        {
            // 确保锁被释放
            _dictLock.Release();
        }
    }

    static async Task Main(string[] args)
    {
        if (true)
        {
            // 创建一组相同的任务
            List<Task> tasks = new List<Task>();
            for (int i = 0; i < _numberOfTasks; i++)
            {
                // 每个任务尝试获取Project1的现有连接或创建新连接
                var t = new Task(async () => { Console.WriteLine(await GetConnection("Project1")); });
                tasks.Add(t);
            }

            // 并行执行这些任务
            Parallel.ForEach<Task>(tasks, (t) => { t.Start(); });
            Task.WaitAll(tasks.ToArray());
            Console.WriteLine("Done");
            Console.Read();
        }
    }
}

核心问题:错误的信号量释放逻辑

你的finally块里无条件调用_dictLock.Release(),这是导致错误的根本原因:

  • 当await _dictLock.WaitAsync(_waitTimeMs)返回false(获取锁超时)时,你并没有成功获取信号量,但finally依然执行Release(),导致信号量的计数超过初始设置的最大值1。
  • 高负载下大量任务超时,多次调用Release()会不断增加信号量计数,最终触发"超过最大计数"的异常。

修复方案

只在成功获取信号量的情况下才释放它,修改逻辑如下:

static async Task<string> GetConnection(string projId)
{
    bool acquiredLock = false;
    try
    {
        // 进入信号量锁,防止同一项目添加多个连接
        acquiredLock = await _dictLock.WaitAsync(_waitTimeMs);
        if (acquiredLock)
        {
            // 尝试从连接池获取连接
            if (_pool.TryGetValue(projId, out bool value))
            {
                if (value == false)
                    return "Exists but not connected yet.";
                else
                    return "Success, exists and connected.";
            }
            // 否则将连接添加到池
            else
            {
                _pool.TryAdd(projId, false);

                // 模拟建立新连接的延迟
                await Task.Delay(2);

                _pool.TryUpdate(projId, true, false);

                return "Created new connection successfully & added to pool.";
            }
        }
        // 报告未能及时获取锁的失败信息
        else
            return "Server busy. Please try again later.";
    }
    catch (Exception ex)
    {
        return "Error " + ex.Message;
    }
    finally
    {
        // 仅在成功获取锁时释放
        if (acquiredLock)
        {
            _dictLock.Release();
        }
    }
}

额外优化建议

  1. 避免混合Task.Start()和Parallel.ForEach:你的Main方法里用Parallel.ForEach启动Task,这是不必要的,直接用Task.WhenAll更符合异步编程规范:
static async Task Main(string[] args)
{
    // 创建一组相同的任务
    List<Task> tasks = new List<Task>();
    for (int i = 0; i < _numberOfTasks; i++)
    {
        // 直接创建异步任务,无需手动Start
        var t = GetConnection("Project1").ContinueWith(t => Console.WriteLine(t.Result));
        tasks.Add(t);
    }

    // 等待所有任务完成
    await Task.WhenAll(tasks);
    Console.WriteLine("Done");
    Console.Read();
}
  1. ConcurrentDictionary的使用优化:既然已经用了信号量保护整个字典操作,其实可以考虑用普通Dictionary替代ConcurrentDictionary,因为信号量已经保证了同一时间只有一个线程操作字典,减少不必要的并发开销。

内容的提问来源于stack exchange,提问作者TheLegendaryCopyCoder

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 07:50:31