C#实现UDP数据接收并同步推送至WebSocket客户端方案问询
解决方案
核心思路
用线程安全的集合管理所有活跃的WebSocket客户端连接,UDP监听线程收到数据后,先将数据存入数据库,再遍历集合将数据推送给每个在线客户端;同时处理客户端连接/断开时的集合更新,确保线程安全。
具体实现步骤
维护线程安全的客户端集合
使用ConcurrentDictionary<string, WebSocket>(用客户端唯一标识做键,方便管理)或者ConcurrentBag<WebSocket>,避免多线程操作时的并发冲突。UDP数据接收与分发逻辑
- UDP线程收到数据后,优先异步执行数据库插入操作(避免阻塞UDP监听流程)
- 遍历WebSocket客户端集合的副本(防止遍历中集合变更引发异常),逐个发送数据;捕获发送异常(比如客户端已断开),并从集合中移除无效连接。
WebSocket连接生命周期管理
- 客户端连接成功时,将其添加到集合
- 客户端断开连接或出现异常时,从集合中移除并释放资源。
代码示例(修改后的服务端核心逻辑)
// 线程安全的WebSocket客户端集合,用Guid生成唯一客户端ID private static readonly ConcurrentDictionary<string, WebSocket> _webSocketClients = new ConcurrentDictionary<string, WebSocket>(); // WebSocket连接处理入口 public async Task HandleWebSocketAsync(HttpContext context) { var webSocket = await context.WebSockets.AcceptWebSocketAsync(); var clientId = Guid.NewGuid().ToString(); // 将客户端加入集合 _webSocketClients.TryAdd(clientId, webSocket); try { var buffer = new byte[1024 * 4]; WebSocketReceiveResult result; do { result = await webSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None); // 如需处理客户端主动发送的消息,可在此添加逻辑 } while (!result.CloseStatus.HasValue); await webSocket.CloseAsync(result.CloseStatus.Value, result.CloseStatusDescription, CancellationToken.None); } catch (Exception ex) { // 记录异常日志(如日志框架Serilog/NLog) } finally { // 断开后移除客户端并释放资源 _webSocketClients.TryRemove(clientId, out var _); webSocket.Dispose(); } } // UDP监听线程逻辑 private async void UdpListenerThread() { using var udpClient = new UdpClient(12345); // 替换为你的UDP监听端口 var remoteIpEndPoint = new IPEndPoint(IPAddress.Any, 0); while (true) { try { // 接收UDP数据 var receiveResult = await udpClient.ReceiveAsync(); var udpData = Encoding.UTF8.GetString(receiveResult.Buffer); // 1. 异步存入数据库,不阻塞UDP监听 await SaveUdpDataToDatabase(udpData); // 2. 推送数据给所有在线WebSocket客户端 var sendBuffer = new ArraySegment<byte>(Encoding.UTF8.GetBytes(udpData)); // 遍历集合副本,避免遍历过程中集合修改引发异常 foreach (var client in _webSocketClients.Values.ToList()) { if (client.State == WebSocketState.Open) { try { await client.SendAsync(sendBuffer, WebSocketMessageType.Text, true, CancellationToken.None); } catch (Exception) { // 发送失败,清理无效客户端 var targetId = _webSocketClients.First(kv => kv.Value == client).Key; _webSocketClients.TryRemove(targetId, out var _); client.Dispose(); } } else { // 客户端已关闭,直接清理 var targetId = _webSocketClients.First(kv => kv.Value == client).Key; _webSocketClients.TryRemove(targetId, out var _); client.Dispose(); } } } catch (Exception ex) { // 记录UDP监听异常日志 } } } // 数据库存储示例方法(替换为你的实际数据库操作) private async Task SaveUdpDataToDatabase(string data) { using var dbContext = new YourDbContext(); // 替换为你的DbContext dbContext.UdpRecords.Add(new UdpRecord { Content = data, ReceiveTime = DateTime.Now }); await dbContext.SaveChangesAsync(); }
关键注意事项
- 线程安全:必须使用线程安全集合,避免多线程添加/移除客户端时的冲突。
- 异步优先级:数据库存储和WebSocket发送都用异步操作,防止阻塞UDP监听线程,保证数据接收的实时性。
- 无效连接清理:发送失败或客户端状态异常时,及时从集合中移除并释放资源,避免内存泄漏。
- 无客户端兼容:即使客户端集合为空,数据库存储逻辑依然执行,完全满足“无连接时仍存数据”的需求。
内容的提问来源于stack exchange,提问作者Joe
相关产品推荐
相关产品推荐

