使用SemaphoreSlim仍出现并发重入,Connected/Disconnected调用异常
问题根源
你遇到的核心问题是**DispatchEventAsync被多次并行执行**,导致Connected事件被连续触发,打破了和Disconnected的交替预期,具体原因如下:
信号量仅保护
ConnectAsync的同步块,未限制DispatchEventAsync的并发执行_connectSemaphore只确保同一时间只有一个线程进入ConnectAsync的try块,但当第一个线程释放信号量后,后续线程会立即进入,此时_dispatch可能还是已完成的任务,于是多个线程会给_dispatch添加多个ContinueWith回调,最终导致多个DispatchEventAsync同时运行,连续触发Connected事件。WebSocket状态未被校验,重复连接未被阻止
ConnectAsync里没有检查_clientSocket是否已经处于连接状态,多次调用ConnectAsync会直接尝试再次连接,进而多次触发Connected事件,而对应的Disconnected可能还未执行(或未执行完成),导致openCloseDifference的差值异常。DispatchEventAsync的生命周期未与WebSocket绑定DispatchEventAsync的执行逻辑没有和WebSocket的实际连接生命周期关联,比如当WebSocket被关闭后,没有机制阻止新的DispatchEventAsync启动,或者确保上一轮的Disconnected执行完成后再触发新的Connected。
修复方案
1. 在ConnectAsync中校验WebSocket状态,避免重复连接
在尝试连接前,先检查_clientSocket的状态,如果已连接则直接返回,同时用锁保护_dispatch的赋值,避免多个回调叠加:
private readonly object _dispatchLock = new object(); public async Task ConnectAsync() { await _connectSemaphore.WaitAsync().ConfigureAwait(false); try { // 检查socket是否已连接,避免重复触发Connected if (_clientSocket?.State == WebSocketState.Open) { return; } await _clientSocket.ConnectAsync().ConfigureAwait(false); // 确保_dispatch的更新是原子性的,避免多个ContinueWith叠加 lock (_dispatchLock) { if (_dispatch.Status == TaskStatus.RanToCompletion) { _dispatch = DispatchEventAsync(); } else { _dispatch = _dispatch.ContinueWith(_ => DispatchEventAsync(), TaskScheduler.Default); } } } finally { _connectSemaphore.Release(); } }
2. 确保DispatchEventAsync与WebSocket生命周期绑定
修改DispatchEventAsync,在执行前检查WebSocket状态,并且在退出时清理socket状态:
private async Task DispatchEventAsync() { try { // 再次确认socket处于连接状态才触发Connected if (_clientSocket?.State != WebSocketState.Open) { return; } Connected(); // 原事件循环逻辑... } catch (Exception exception) { } finally { try { Disconnected(); // 清理socket状态,避免后续重复连接判断出错 if (_clientSocket?.State != WebSocketState.Closed) { await _clientSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None).ConfigureAwait(false); } _clientSocket = null; } catch (Exception ex) { } } }
3. 优化连接状态跟踪逻辑
用线程安全的状态枚举替代openCloseDifference,避免简单增减带来的并发问题:
private enum ConnectionState { Disconnected, Connecting, Connected } private readonly System.Threading.Atomic<ConnectionState> _connectionState = new(System.Threading.Atomic<ConnectionState>(ConnectionState.Disconnected));
在Connected和Disconnected事件触发前,先原子性切换状态,确保状态变更的一致性。
4. 规范CloseAsync的实现
确保CloseAsync能正确终止DispatchEventAsync的事件循环,并且更新连接状态:
public async Task CloseAsync() { await _connectSemaphore.WaitAsync().ConfigureAwait(false); try { if (_clientSocket?.State == WebSocketState.Open) { await _clientSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None).ConfigureAwait(false); // 给events通道写入终止信号,让DispatchEventAsync的循环退出 await events.Writer.WriteAsync(new TerminateEvent()).ConfigureAwait(false); } } finally { _connectSemaphore.Release(); } }
额外建议
- 避免在压力测试中无限制地并发调用
ConnectAsync和CloseAsync,可以用SemaphoreSlim限制同时执行的任务数量,减少不必要的并发冲突。 - 给
events通道添加明确的终止机制,确保DispatchEventAsync的循环能正常退出,避免僵尸任务占用资源。
内容的提问来源于stack exchange,提问作者jeiea

