C# TCP服务器接收客户端消息延迟问题排查求助
TCP服务器消息延迟问题排查求助
我在C#应用中实现了基于TCP Listener的TCP服务器,应用同时包含MQTT和Kafka模块。目前遇到一个棘手问题:TCP客户端发送的消息有时会延迟数分钟到2小时才到达服务器,延迟结束后所有积压消息会一次性全部送达。
已做排查
- 通过Wireshark抓包验证(拓扑:TCP Client ---> Wireshark ---> TCP Server),确认消息在客户端发送时已到达抓包节点,排除客户端和网络链路问题,判断问题出在TCP服务器端
- 已按照建议将TCP服务器改为后台运行:
tcpServer = new TcpServer(); _ = Task.Run(tcpServer.Run); - 本地测试无法复现该问题
相关代码
1. 启动TCP监听器
public async void Run() { try { if (!IsRunning) { IsRunning = true; clientTokenSource = new CancellationTokenSource(); clientListener = new TcpListener(Ip, ServerPort); clientListener.Start(); while (true) { try { TcpClient RemoteClient = await clientListener.AcceptTcpClientAsync(); IPEndPoint clientIP = RemoteClient.Client.RemoteEndPoint as IPEndPoint; _ = InitializePmsClient(RemoteClient); } catch (Exception _ex) { Logger.Error(_ex); break; } } // Manually close all connected clients for (int i = clientCollection.Count - 1; i >= 0; i--) { TcpClient tcpClient = clientCollection[i].TcpClient; IPEndPoint clientIP = tcpClient.Client.RemoteEndPoint as IPEndPoint; tcpClient.Close(); clientCollection[i].Connected = false; Logger.Info("Disconnected to client"); } // Wait for all of the clients to finish closing await Task.WhenAll(clientTasks); clientTasks.Clear(); clientTokenSource.Dispose(); clientCollection.Clear(); } else { clientTokenSource.Cancel(); clientListener.Stop(); IsRunning = false; } } catch (Exception ex) { Logger.Error(ex); } }
2. 初始化TCP客户端
private PmsClient InitializePmsClient(TcpClient remoteClient) { PmsClient client = new PmsClient(Logger) { TcpClient = remoteClient }; client.TcpClient.NoDelay = true; client.Connected = true; client.KeepAlivePeriod = KeepAlivePeriod; client.ProcessDataTask = ProcessClientAsync(client, clientTokenSource.Token); // Process when received message client.SendDataTask = TaskSendMessage(client); // Dequeue message from queue to client client.SendKeepAlive(); // Send keep alive message periodically clientCollection.Add(client); Logger.Info("Client connected"); return client; }
3. 处理客户端消息
private async Task ProcessClientAsync(PmsClient client, CancellationToken token) { int max_read_len = 4096; try { // Begin reading from the client's data stream using (NetworkStream stream = client.TcpClient.GetStream()) { int read = 1; while (read > 0) { byte[] buffer = new byte[client.TcpClient.ReceiveBufferSize]; read = await stream.ReadAsync(buffer, 0, max_read_len, token); client.ProcessData(buffer, read); // Format data from buffer } } } catch (Exception ex) { Logger.Error(ex); } finally { try { client.TcpClient.Close(); clientTasks.Remove(client.ProcessDataTask); clientTasks.Remove(client.SendDataTask); clientCollection.Remove(client); Logger.Info("Client disconnected"); } catch (Exception exception) { Logger.Error(exception); } } }
4. 处理接收到的数据细节
public void ProcessData(byte[] data, int read) { PmsMessage Message = new PmsMessage(); Booking booking; try { if (read > 0) { Received.AddRange(data.Take(read)); Logger.Info("Received data", Received); // <---------------------- Logging data when received data from client // Read from buffer until message receive all message bool CompleteMessage = Message.IsCompleteMessage(ref Received, out result); while (CompleteMessage) { // Trigger event handle message received logic OnMessageTrigger(new MessageEventArgs(result)); CompleteMessage = Message.IsCompleteMessage(ref Received, out result); } } } catch (Exception ex) { Logger.Error(ex); } }
排查思路与解决方案建议
- 检查线程池资源:监控延迟发生时
ThreadPool.GetAvailableThreads()的数值,确认是否因MQTT/Kafka模块耗尽线程池,导致TCP接收的异步回调无法及时执行 - 排查
ProcessData中的阻塞操作:OnMessageTrigger事件处理如果包含同步阻塞逻辑(如Kafka同步发送、数据库同步读写),会卡住接收流程,建议将事件处理异步化 - 优化TCP缓冲区配置:手动设置合理的固定接收缓冲区大小(如8192),同时检查系统层面的TCP缓冲区参数(Windows注册表
TcpWindowSize、Linuxnet.ipv4.tcp_rmem) - 优化日志性能:
Logger.Info("Received data", Received)若为同步写入,高并发下会阻塞,建议改用异步日志或降低该日志级别 - 确保集合操作线程安全:
clientTasks和clientCollection的增删操作未加锁,多客户端场景下可能引发线程安全问题,需添加锁保护 - 验证CancellationToken传递:确认令牌是否正确传递到所有异步操作,避免因令牌未触发导致的流程阻塞
- 监控系统资源指标:延迟发生时检查服务器CPU、内存、磁盘IO、网络连接数,以及TCP连接状态(
netstat/ss命令),排查资源耗尽或连接异常情况
内容的提问来源于stack exchange,提问作者simpsons3
相关产品推荐
相关产品推荐

