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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 05:10:21