如何为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
相关产品推荐
相关产品推荐

