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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 23:10:26