ASP.NET MVC WebSocket保活与业务指令响应冲突问题咨询
WebSocket保活与业务消息共存的解决方案
核心问题分析
你的问题本质是WebSocket连接的消息接收逻辑被ScheduleService独占,保活的循环接收拦截了所有消息却未做分发,导致控制器无法获取业务响应,仅能偶尔收到未被处理的echo回复。
具体实现方案
1. 重构ScheduleService,实现消息分发机制
将原有的ReceiveMessagesAsync改造成同时处理保活和业务消息的逻辑,通过事件委托把业务消息转发给控制器,避免消息被独占。
示例代码:
public class ScheduleService { private readonly WebSocket _webSocket; private CancellationTokenSource _cts; // 定义业务消息事件,供控制器订阅 public event Action<string> OnBusinessMessageReceived; public ScheduleService(WebSocket webSocket) { _webSocket = webSocket; _cts = new CancellationTokenSource(); } public async Task StartKeepAliveAndMessageHandling() { // 启动60秒间隔的保活任务 var keepAliveTask = Task.Run(async () => { while (!_cts.Token.IsCancellationRequested) { await SendEchoAsync(); await Task.Delay(60000, _cts.Token); } }, _cts.Token); // 启动单线程消息接收循环 var receiveTask = ReceiveMessagesLoopAsync(); await Task.WhenAll(keepAliveTask, receiveTask); } private async Task SendEchoAsync() { var echoMsg = JsonSerializer.Serialize(new { type = "echo" }); var buffer = Encoding.UTF8.GetBytes(echoMsg); await _webSocket.SendAsync(new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None); } private async Task ReceiveMessagesLoopAsync() { var buffer = new byte[1024 * 4]; while (!_cts.Token.IsCancellationRequested && _webSocket.State == WebSocketState.Open) { var result = await _webSocket.ReceiveAsync(new ArraySegment<byte>(buffer), _cts.Token); if (result.MessageType == WebSocketMessageType.Close) { await _webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closed by client", CancellationToken.None); break; } var message = Encoding.UTF8.GetString(buffer, 0, result.Count); var msgObj = JsonSerializer.Deserialize<dynamic>(message); // 过滤echo回复,仅转发业务消息 if (msgObj?.type != "echo") { OnBusinessMessageReceived?.Invoke(message); } } } public void Stop() { _cts.Cancel(); } }
2. 控制器订阅业务消息并处理
在控制器的WebSocket连接初始化逻辑中,实例化ScheduleService并订阅消息事件,确保业务响应能被控制器接收处理。
示例代码:
[ApiController] [Route("api/schedules")] public class ScheduleController : ControllerBase { [HttpGet("connect")] public async Task Connect() { if (!HttpContext.WebSockets.IsWebSocketRequest) { HttpContext.Response.StatusCode = StatusCodes.Status400BadRequest; return; } using var webSocket = await HttpContext.WebSockets.AcceptWebSocketAsync(); var scheduleService = new ScheduleService(webSocket); // 订阅业务消息回调 scheduleService.OnBusinessMessageReceived += HandleBusinessMessage; try { // 先发送登录消息,再启动保活和消息处理 await SendLoginMessage(webSocket); await scheduleService.StartKeepAliveAndMessageHandling(); } finally { scheduleService.Stop(); scheduleService.OnBusinessMessageReceived -= HandleBusinessMessage; } } private void HandleBusinessMessage(string message) { // 解析并处理业务消息,比如getSchedules的响应 var scheduleResponse = JsonSerializer.Deserialize<ScheduleResponse>(message); // 此处可根据需求将数据返回前端或存储 } private async Task SendLoginMessage(WebSocket webSocket) { var loginMsg = JsonSerializer.Serialize(new { type = "login", userId = "your-user-id" }); var buffer = Encoding.UTF8.GetBytes(loginMsg); await webSocket.SendAsync(new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None); } // 按需调用发送getSchedules指令 public async Task SendGetSchedules(WebSocket webSocket, DateTime targetDate) { var getSchedMsg = JsonSerializer.Serialize(new { type = "getSchedules", date = targetDate.ToString("yyyy-MM-dd") }); var buffer = Encoding.UTF8.GetBytes(getSchedMsg); await webSocket.SendAsync(new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None); } }
3. 关键注意事项
- 同一个WebSocket连接只能有一个活跃的
ReceiveAsync循环,避免多线程抢占导致消息混乱。 - 增加异常捕获逻辑,处理WebSocket断开、网络异常等场景,及时取消任务避免资源泄漏。
- 确保消息序列化/反序列化格式与服务器完全匹配,避免解析失败。
内容的提问来源于stack exchange,提问作者Ramees Shahzad
相关产品推荐
相关产品推荐

