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

使用SemaphoreSlim仍出现并发重入,Connected/Disconnected调用异常

问题分析与解决方案

问题根源

你遇到的核心问题是**DispatchEventAsync被多次并行执行**,导致Connected事件被连续触发,打破了和Disconnected的交替预期,具体原因如下:

  1. 信号量仅保护ConnectAsync的同步块,未限制DispatchEventAsync的并发执行
    _connectSemaphore只确保同一时间只有一个线程进入ConnectAsync的try块,但当第一个线程释放信号量后,后续线程会立即进入,此时_dispatch可能还是已完成的任务,于是多个线程会给_dispatch添加多个ContinueWith回调,最终导致多个DispatchEventAsync同时运行,连续触发Connected事件。

  2. WebSocket状态未被校验,重复连接未被阻止
    ConnectAsync里没有检查_clientSocket是否已经处于连接状态,多次调用ConnectAsync会直接尝试再次连接,进而多次触发Connected事件,而对应的Disconnected可能还未执行(或未执行完成),导致openCloseDifference的差值异常。

  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:45:34