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

如何为WebSocket添加代理并解决多用户下服务器垃圾消息问题

解决多用户场景下.NET WebSocket客户端的「垃圾消息」问题:服务器端Proxy方案

原来的代码在单用户场景下直接连接第三方WebSocket服务没问题,但多用户时每个用户都创建独立连接,会导致第三方服务推送的消息被重复处理,加上服务器端连接管理混乱,就出现了「垃圾消息」。解决思路是在服务器端实现Proxy代理类,只维护一个到第三方WebSocket的长连接,由Proxy统一管理所有用户的订阅,把收到的行情消息广播给所有用户,用户端只和自己服务器的WebSocket交互,不再直接连接第三方。


1. 实现WebSocket Proxy核心类

这个类负责维护与第三方的连接、管理用户订阅、广播消息,用单例确保只有一个第三方连接:

using System.Collections.Concurrent;
using System.Net.WebSockets;
using Websocket.Client;

namespace TraderMadeWebSocketProxy
{
    public class MarketDataWebSocketProxy
    {
        // 单例实例,确保全局唯一
        private static readonly Lazy<MarketDataWebSocketProxy> _instance = new(() => new MarketDataWebSocketProxy());
        public static MarketDataWebSocketProxy Instance => _instance.Value;

        // 第三方WebSocket客户端实例
        private WebsocketClient? _thirdPartyClient;
        // 线程安全的用户订阅集合,存储用户ID和对应的WebSocket连接
        private readonly ConcurrentDictionary<string, WebSocket> _subscribedUsers = new();
        private readonly string _streamingApiKey = "streaming_api_key";
        private readonly Uri _thirdPartyUrl = new("wss://marketdata.tradermade.com/feedadv");

        // 私有构造函数,禁止外部实例化
        private MarketDataWebSocketProxy() { }

        // 初始化第三方WebSocket连接
        public async Task InitializeAsync()
        {
            if (_thirdPartyClient != null && _thirdPartyClient.IsRunning)
                return;

            _thirdPartyClient = new WebsocketClient(_thirdPartyUrl);
            _thirdPartyClient.ReconnectTimeout = TimeSpan.FromSeconds(30);

            // 处理重连事件,重连后重新发送订阅请求
            _thirdPartyClient.ReconnectionHappened.Subscribe(info =>
            {
                Console.WriteLine($"第三方WebSocket重连,类型:{info.Type}");
                SendSubscriptionRequest();
            });

            // 接收第三方消息并广播给所有订阅用户
            _thirdPartyClient.MessageReceived.Subscribe(async msg =>
            {
                Console.WriteLine($"收到第三方消息:{msg}");
                if (msg.ToString().ToLower() == "connected")
                {
                    SendSubscriptionRequest();
                    return;
                }

                await BroadcastMessageToUsers(msg.ToString());
            });

            await _thirdPartyClient.StartOrFail();
        }

        // 发送订阅请求到第三方服务
        private void SendSubscriptionRequest()
        {
            if (_thirdPartyClient == null || !_thirdPartyClient.IsRunning)
                return;

            string subscribeData = $"{{\"userKey\":\"{_streamingApiKey}\", \"symbol\":\"EURUSD,GBPUSD,USDJPY\"}}";
            _thirdPartyClient.Send(subscribeData);
            Console.WriteLine("已发送订阅请求到第三方服务");
        }

        // 添加用户到订阅列表
        public bool AddUserSubscription(string userId, WebSocket userWebSocket)
        {
            return _subscribedUsers.TryAdd(userId, userWebSocket);
        }

        // 移除用户订阅
        public bool RemoveUserSubscription(string userId)
        {
            return _subscribedUsers.TryRemove(userId, out _);
        }

        // 广播消息给所有订阅用户,自动清理失效连接
        private async Task BroadcastMessageToUsers(string message)
        {
            var buffer = new ArraySegment<byte>(System.Text.Encoding.UTF8.GetBytes(message));
            var invalidUsers = new List<string>();

            foreach (var (userId, socket) in _subscribedUsers)
            {
                if (socket.State != WebSocketState.Open)
                {
                    invalidUsers.Add(userId);
                    continue;
                }

                try
                {
                    await socket.SendAsync(buffer, WebSocketMessageType.Text, true, CancellationToken.None);
                }
                catch (Exception ex)
                {
                    Console.WriteLine($"给用户{userId}发送消息失败:{ex.Message}");
                    invalidUsers.Add(userId);
                }
            }

            // 清理失效的用户连接
            foreach (var userId in invalidUsers)
            {
                RemoveUserSubscription(userId);
            }
        }

        // 停止Proxy服务
        public async Task StopAsync()
        {
            if (_thirdPartyClient != null)
            {
                await _thirdPartyClient.StopOrFail();
                _thirdPartyClient.Dispose();
            }
        }
    }
}

2. 服务器端处理用户WebSocket连接(ASP.NET Core示例)

通过中间件处理用户的WebSocket连接请求,与Proxy交互完成订阅/取消订阅:

using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Http;
using System.Net.WebSockets;

namespace TraderMadeWebSocketProxy
{
    public static class WebSocketMiddlewareExtensions
    {
        public static IApplicationBuilder UseMarketDataWebSocket(this IApplicationBuilder app)
        {
            return app.UseMiddleware<MarketDataWebSocketMiddleware>();
        }
    }

    public class MarketDataWebSocketMiddleware
    {
        private readonly RequestDelegate _next;

        public MarketDataWebSocketMiddleware(RequestDelegate next)
        {
            _next = next;
        }

        public async Task InvokeAsync(HttpContext context)
        {
            // 只处理WebSocket请求
            if (!context.WebSockets.IsWebSocketRequest)
            {
                context.Response.StatusCode = StatusCodes.Status400BadRequest;
                return;
            }

            // 这里的用户ID可替换为实际认证信息(如JWT、Session),示例用随机ID
            var userId = Guid.NewGuid().ToString();
            using var userWebSocket = await context.WebSockets.AcceptWebSocketAsync();

            // 初始化Proxy(确保全局只初始化一次)
            await MarketDataWebSocketProxy.Instance.InitializeAsync();
            // 添加用户到订阅列表
            MarketDataWebSocketProxy.Instance.AddUserSubscription(userId, userWebSocket);

            Console.WriteLine($"用户{userId}已连接并订阅行情");

            try
            {
                // 监听用户消息(可扩展处理用户自定义订阅请求,比如修改订阅品种)
                var buffer = new byte[1024 * 4];
                WebSocketReceiveResult result;
                do
                {
                    result = await userWebSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None);
                    // 若需处理用户自定义请求,可解析buffer中的内容并修改Proxy订阅逻辑
                } while (!result.CloseStatus.HasValue);

                // 用户主动断开连接
                await userWebSocket.CloseAsync(result.CloseStatus.Value, result.CloseStatusDescription, CancellationToken.None);
            }
            catch (Exception ex)
            {
                Console.WriteLine($"用户{userId}连接异常:{ex.Message}");
            }
            finally
            {
                // 移除用户订阅
                MarketDataWebSocketProxy.Instance.RemoveUserSubscription(userId);
                Console.WriteLine($"用户{userId}已取消订阅并断开连接");
            }
        }
    }
}

3. 在ASP.NET Core中注册中间件

在Startup类中启用WebSocket支持并注册自定义中间件:

using Microsoft.AspNetCore.Builder;
using Microsoft.Extensions.DependencyInjection;

namespace TraderMadeWebSocketProxy
{
    public class Startup
    {
        public void ConfigureServices(IServiceCollection services)
        {
            // 注册其他必要服务
        }

        public void Configure(IApplicationBuilder app)
        {
            // 启用WebSocket支持
            app.UseWebSockets();
            // 注册行情WebSocket中间件
            app.UseMarketDataWebSocket();

            // 其他中间件配置(如路由、静态文件等)
        }
    }
}

方案优势

  • 单一连接:服务器端仅维护一个到第三方的长连接,避免多用户重复连接导致的消息混乱。
  • 高效管理:线程安全集合存储用户连接,自动清理失效连接,避免资源泄漏。
  • 消息一致:第三方推送的消息统一广播给所有用户,确保消息无重复、无混乱。
  • 可扩展性:可快速扩展用户自定义订阅逻辑(如用户指定订阅品种)。

内容的提问来源于stack exchange,提问作者Rami

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 09:07:43