如何高效实现WebSocket响应等待?替代Task.Delay的优化方案
优化WebSocket响应获取方案(替代轮询)
你当前用Task.Delay轮询字典的方案存在延迟高、资源浪费的问题,用TaskCompletionSource<T>可以实现响应到达时立即唤醒等待任务,完全替代轮询逻辑,是更高效的实现方式。
核心思路
用线程安全的字典存储每个请求ID对应的等待任务(TaskCompletionSource<string>),当WebSocket收到响应时,直接找到对应请求的任务并标记完成,等待的代码就能立即拿到结果,无需轮询。
完整实现代码
1. 替换存储结构
用ConcurrentDictionary替代原有的普通字典,保证多线程下的安全性(发送请求和接收响应通常在不同线程执行):
private readonly ConcurrentDictionary<string, TaskCompletionSource<string>> _pendingRequests = new(); private readonly WebSocket _socket; // 你的WebSocket实例
2. 发送请求并等待响应的方法
public async Task<string> SendRequestAndWaitForResponse(byte[] binaryMessage, string requestId, CancellationToken cancellationToken = default) { // 创建等待任务并加入字典 var tcs = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously); if (!_pendingRequests.TryAdd(requestId, tcs)) { throw new InvalidOperationException($"请求ID {requestId} 重复,请勿重复发送"); } try { // 发送WebSocket消息 await _socket.SendAsync( new ArraySegment<byte>(binaryMessage, 0, binaryMessage.Length), WebSocketMessageType.Text, endOfMessage: true, cancellationToken); // 设置60秒超时,同时支持外部取消 using var timeoutToken = new CancellationTokenSource(TimeSpan.FromSeconds(60)); using var linkedToken = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, timeoutToken.Token); return await tcs.Task.WaitAsync(linkedToken.Token); } catch (OperationCanceledException) { // 超时或取消时清理任务 _pendingRequests.TryRemove(requestId, out _); return null; // 也可根据业务需求抛出异常 } finally { // 无论成功失败,都清理字典中的任务,防止内存泄漏 _pendingRequests.TryRemove(requestId, out _); } }
3. WebSocket接收消息的处理逻辑
修改你原有的接收循环,收到响应时直接触发对应任务的完成:
// 假设这是你的WebSocket接收循环方法 private async Task ReceiveLoop(CancellationToken cancellationToken) { var buffer = new byte[1024 * 4]; while (!cancellationToken.IsCancellationRequested && _socket.State == WebSocketState.Open) { var receiveResult = await _socket.ReceiveAsync(new ArraySegment<byte>(buffer), cancellationToken); if (receiveResult.MessageType == WebSocketMessageType.Close) { await _socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "连接关闭", cancellationToken); break; } // 解析响应消息,提取requestId和响应内容(根据你的实际消息格式调整) string responseJson = Encoding.UTF8.GetString(buffer, 0, receiveResult.Count); var socketResult = JsonSerializer.Deserialize<SocketResult>(responseJson); // 你的SocketResult类型 // 找到对应请求的任务并完成它 if (_pendingRequests.TryRemove(socketResult.RequestId, out var tcs)) { tcs.SetResult(socketResult.Json); } else { // 处理超时后迟来的响应,可记录日志或直接忽略 } } }
对比原方案的优势
- 零延迟:响应到达瞬间唤醒等待任务,无需轮询等待
- 资源高效:避免了定时
Task.Delay的资源消耗,以及轮询时的重复字典查询 - 线程安全:用
ConcurrentDictionary处理并发场景,不会出现字典操作冲突 - 自动清理:通过
finally块和超时机制自动移除过期任务,防止内存泄漏
注意事项
- 确保消息解析能正确提取
requestId,这是关联请求与响应的核心 - 超时时间可根据业务需求调整,也可通过方法参数动态传入
- 处理WebSocket异常、关闭等场景时,需遍历并取消所有未完成的任务,避免内存泄漏
TaskCreationOptions.RunContinuationsAsynchronously可避免同步延续导致的线程阻塞问题
内容的提问来源于stack exchange,提问作者mayi
相关产品推荐
相关产品推荐

