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

多线程流读写与保活机制实现:TCP客户端开发技术问询

咱们先理清楚你当前代码里的几个核心问题,然后给你一套能同时处理业务消息发送、响应监听和后台保活的可行方案:

核心问题分析

你的需求是客户端要同时做三件事:接收用户输入发业务消息、监听服务器所有响应、每15秒自动发保活,但现有代码存在这些问题:

  • 不必要创建多个TcpClient连接,客户端和服务器保持一个长连接就足够,多连接反而会增加复杂度
  • Send方法被用户输入完全阻塞,根本没法同时处理定时保活逻辑
  • 每次循环都新建StreamWriter/StreamReader,会导致流状态混乱,还浪费资源
  • Read方法依赖stream.DataAvailable判断,很容易错过服务器发送的消息(比如消息到来时刚好不在循环判断的瞬间)
  • 没有独立的保活线程,当前的Thread.Sleep(23000)完全不符合15秒保活的要求

解决方案实现

下面是调整后的完整代码,采用单长连接+三线程分工的架构,完美满足你的需求:

using System;
using System.Collections.Generic;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;

public class TcpClientWithKeepAlive
{
    private static TcpClient _tcpClient;
    private static NetworkStream _networkStream;
    private static StreamWriter _streamWriter;
    private static StreamReader _streamReader;
    
    // 连接同步信号量
    private static ManualResetEvent _connectDone = new ManualResetEvent(false);
    // 全局取消令牌,用于优雅终止所有线程
    private static CancellationTokenSource _cts = new CancellationTokenSource();

    public static void Main()
    {
        try
        {
            // 建立单个长连接
            Connect("你的服务器IP", 你的端口号);
            _connectDone.WaitOne();
            
            // 初始化流对象,只创建一次,避免重复实例化
            _networkStream = _tcpClient.GetStream();
            _streamWriter = new StreamWriter(_networkStream, Encoding.UTF8) { AutoFlush = false };
            _streamReader = new StreamReader(_networkStream, Encoding.UTF8);
            
            // 启动三个分工线程
            var readThread = new Thread(() => ReadServerResponses(_cts.Token));
            var businessSendThread = new Thread(() => SendBusinessMessages(_cts.Token));
            var keepAliveThread = new Thread(() => SendKeepAliveMessages(_cts.Token));
            
            readThread.Start();
            businessSendThread.Start();
            keepAliveThread.Start();
            
            Console.WriteLine("客户端启动完成,按任意键退出...");
            Console.ReadKey();
            
            // 优雅终止所有线程和连接
            _cts.Cancel();
            _tcpClient.Close();
        }
        catch (Exception ex)
        {
            Console.WriteLine($"客户端异常: {ex.Message}");
        }
    }

    private static void Connect(string ip, int port)
    {
        try
        {
            _tcpClient = new TcpClient();
            _tcpClient.BeginConnect(ip, port, ar =>
            {
                try
                {
                    _tcpClient.EndConnect(ar);
                    Console.WriteLine("成功连接到服务器");
                    _connectDone.Set();
                }
                catch (Exception ex)
                {
                    Console.WriteLine($"连接失败: {ex.Message}");
                }
            }, null);
        }
        catch (Exception ex)
        {
            Console.WriteLine($"连接异常: {ex.Message}");
        }
    }

    // 业务消息发送线程:专门处理用户输入,发送业务消息
    private static void SendBusinessMessages(CancellationToken token)
    {
        while (!token.IsCancellationRequested && _tcpClient.Connected)
        {
            try
            {
                Console.WriteLine("\n请输入消息索引(整数):");
                Console.Write("输入: ");
                var input = Console.ReadLine();
                
                // 校验输入合法性
                if (!int.TryParse(input, out int msgIndex))
                {
                    Console.WriteLine("请输入有效的整数!");
                    continue;
                }
                
                if (msgIndex < 0 || msgIndex >= Messages.Count)
                {
                    Console.WriteLine("消息索引超出范围!");
                    continue;
                }
                
                // 加锁确保多线程写流的线程安全
                lock (_networkStream)
                {
                    var message = DPLHeader(Messages[msgIndex]);
                    _streamWriter.WriteLine(message);
                    _streamWriter.Flush();
                    Console.WriteLine($"已发送业务消息: {message}");
                }
                
                // 可选:避免用户频繁输入,按需调整
                Thread.Sleep(1000);
            }
            catch (Exception ex)
            {
                Console.WriteLine($"发送业务消息异常: {ex.Message}");
                break;
            }
        }
    }

    // 响应监听线程:持续读取服务器所有响应(业务+保活)
    private static void ReadServerResponses(CancellationToken token)
    {
        while (!token.IsCancellationRequested && _tcpClient.Connected)
        {
            try
            {
                // 用异步读取替代DataAvailable判断,避免错过消息
                var response = await _streamReader.ReadLineAsync();
                
                if (response == null)
                {
                    Console.WriteLine("服务器已断开连接");
                    _cts.Cancel();
                    break;
                }
                
                Console.WriteLine($"收到服务器响应: {response}");
                
                // 解析并处理响应
                var decodedMsg = DecodingMsg(response);
                var resultList = Decode(decodedMsg);
                SendToText(resultList);
            }
            catch (Exception ex)
            {
                Console.WriteLine($"读取响应异常: {ex.Message}");
                break;
            }
        }
    }

    // 保活线程:每15秒自动发送保活消息
    private static void SendKeepAliveMessages(CancellationToken token)
    {
        // 按服务器要求的格式构造保活消息
        var keepAliveMsg = DPLHeader("KEEP_ALIVE");
        
        while (!token.IsCancellationRequested && _tcpClient.Connected)
        {
            try
            {
                lock (_networkStream)
                {
                    _streamWriter.WriteLine(keepAliveMsg);
                    _streamWriter.Flush();
                    Console.WriteLine($"已发送保活消息: {keepAliveMsg}");
                }
                
                // 等待15秒,同时监听取消信号
                if (token.WaitHandle.WaitOne(15000))
                {
                    break;
                }
            }
            catch (Exception ex)
            {
                Console.WriteLine($"发送保活消息异常: {ex.Message}");
                break;
            }
        }
    }

    // 以下是你原有代码中的业务方法,直接保留即可
    private static string DPLHeader(string msg)
    {
        // 你的消息头封装逻辑
        return msg;
    }

    private static string DecodingMsg(string response)
    {
        // 你的消息解码逻辑
        return response;
    }

    private static List<string> Decode(string decodedMsg)
    {
        // 你的消息解析逻辑
        return new List<string> { decodedMsg };
    }

    private static void SendToText(List<string> list)
    {
        // 你的结果写入文本逻辑
        foreach (var item in list)
        {
            Console.WriteLine($"解析结果: {item}");
        }
    }

    // 假设你的业务消息列表
    private static List<string> Messages = new List<string> { "业务消息1", "业务消息2", "业务消息3" };
}

关键改进点说明

  1. 单长连接复用:所有通信通过一个TcpClient完成,符合TCP长连接的设计规范
  2. 三线程分工:
    • 业务发送线程:专注处理用户输入和业务消息发送
    • 响应监听线程:异步读取服务器所有响应,不会错过消息
    • 保活线程:独立定时发送保活消息,不受其他逻辑阻塞
  3. 流对象复用:只初始化一次StreamWriter/StreamReader,避免流状态混乱
  4. 线程安全:所有写流操作加锁,防止多线程同时操作流导致的异常
  5. 优雅退出:用CancellationTokenSource统一管理线程终止,避免资源泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:39:15