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

C#中使用DataContractJsonSerializer流式传输对象,保持Socket连接的问题

问题:保持Socket打开时,DataContractJsonSerializer无法在接收端即时获取数据

我尝试在C#中使用DataContractJsonSerializer类通过网络发送特定类的对象,但接收方只有在发送方关闭Socket时才能收到数据。现在需要保持Socket处于打开状态以便后续发送更多数据,该如何处理?

待传输的类

[DataContract]
public class Message
{
    [DataMember]
    public DateTime datetime;

    [DataMember]
    public string contents;

    public Message()
    {
        datetime = DateTime.Now;
        contents = "Hello World";
    }

    public void WriteJson(Stream stream)
    {
        var serializer = new DataContractJsonSerializer(typeof(Message));
        serializer.WriteObject(stream, this);
    }

    public static Message ReadJson(Stream stream)
    {
        var serializer = new DataContractJsonSerializer(typeof(Message));
        return (Message)serializer.ReadObject(stream);
    }

    public override string ToString()
    {
        using (MemoryStream stream = new MemoryStream())
        {
            WriteJson(stream);
            return Encoding.UTF8.GetString(stream.ToArray());
        }
    }
}

客户端代码

public static void Main(string[] args)
{
    IPAddress adr = IPAddress.Loopback;
    IPEndPoint ep = new IPEndPoint(adr, 4711);
    Socket clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
    while (true)
    {
        try
        {
            Console.WriteLine("Trying to connect ...");
            clientSocket.Connect(ep);
            NetworkStream clientSocketStream = new NetworkStream(clientSocket);

            //receive data
            Message m = Message.ReadJson(clientSocketStream);
            Console.WriteLine(m);
            clientSocket.Shutdown(SocketShutdown.Both);
            clientSocket.Close();
            Console.ReadLine();
            return;
        }
        catch (SocketException se)
        {
            Console.WriteLine(se.Message);
        }
    }
}

服务端代码

static void Main(string[] args)
{
    const int PORT = 4711;

    IPAddress adr = IPAddress.Any;
    IPEndPoint ep = new IPEndPoint(adr, PORT);
    Socket serverSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
    serverSocket.Bind(ep);
    serverSocket.Listen(20);
    while (true)
    {
        Console.WriteLine("Waiting for connection request on port {0} ...", PORT);
        Socket clientSocket = serverSocket.Accept();
        NetworkStream clientSocketStream = new NetworkStream(clientSocket);

        Console.WriteLine("Hello is sent to:  " + clientSocket.RemoteEndPoint);
        Message m = new Message();
        m.WriteJson(clientSocketStream);
        //clientSocketStream.Flush();

        Console.Write("Wait for key input ...");
        Console.ReadKey();
        //Message is received at the client only when the follwing code is execured
        clientSocket.Shutdown(SocketShutdown.Both);
        clientSocket.Close();
    }
    serverSocket.Close();
}
解决方案

问题原因

DataContractJsonSerializer.ReadObject()方法会持续等待输入流的结束信号(只有当Socket关闭时,流才会触发结束)。而TCP是无消息边界的流协议,如果不明确告知接收方单条消息的长度,接收端无法判断当前消息是否传输完成,只能一直等待流关闭。

解决思路:基于长度的帧分割

在每次发送消息前,先发送消息的字节长度(用固定字节数的整数表示,比如4字节的int),接收端先读取长度,再根据长度读取对应字节数的数据,即可在不关闭Socket的情况下完成单条消息的接收。

修改后的代码

1. 扩展Message类,增加序列化字节数组的方法

[DataContract]
public class Message
{
    // 保留原有成员和方法...

    // 新增:将对象序列化为字节数组
    public byte[] ToJsonBytes()
    {
        using (MemoryStream stream = new MemoryStream())
        {
            WriteJson(stream);
            return stream.ToArray();
        }
    }
}

2. 修改服务端发送逻辑

static void Main(string[] args)
{
    const int PORT = 4711;

    IPAddress adr = IPAddress.Any;
    IPEndPoint ep = new IPEndPoint(adr, PORT);
    Socket serverSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
    serverSocket.Bind(ep);
    serverSocket.Listen(20);
    while (true)
    {
        Console.WriteLine("Waiting for connection request on port {0} ...", PORT);
        Socket clientSocket = serverSocket.Accept();
        NetworkStream clientSocketStream = new NetworkStream(clientSocket);

        Console.WriteLine("Hello is sent to:  " + clientSocket.RemoteEndPoint);
        Message m = new Message();
        byte[] jsonBytes = m.ToJsonBytes();
        
        // 步骤1:先发送消息的字节长度(4字节int)
        byte[] lengthBytes = BitConverter.GetBytes(jsonBytes.Length);
        clientSocketStream.Write(lengthBytes, 0, lengthBytes.Length);
        
        // 步骤2:发送实际的消息数据
        clientSocketStream.Write(jsonBytes, 0, jsonBytes.Length);
        clientSocketStream.Flush(); // 强制将缓冲区数据发送到网络

        Console.Write("Socket保持打开,按任意键可发送下一条消息或关闭连接...");
        Console.ReadKey();
        
        // 此处不再立即关闭Socket,可添加后续发送逻辑
        // clientSocket.Shutdown(SocketShutdown.Both);
        // clientSocket.Close();
    }
    serverSocket.Close();
}

3. 修改客户端接收逻辑

public static void Main(string[] args)
{
    IPAddress adr = IPAddress.Loopback;
    IPEndPoint ep = new IPEndPoint(adr, 4711);
    Socket clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
    while (true)
    {
        try
        {
            Console.WriteLine("Trying to connect ...");
            clientSocket.Connect(ep);
            NetworkStream clientSocketStream = new NetworkStream(clientSocket);

            // 步骤1:读取消息长度(4字节int)
            byte[] lengthBytes = new byte[4];
            int readLength = clientSocketStream.Read(lengthBytes, 0, 4);
            if (readLength != 4)
            {
                Console.WriteLine("读取消息长度失败");
                break;
            }
            int dataLength = BitConverter.ToInt32(lengthBytes, 0);

            // 步骤2:读取对应长度的消息数据(循环读取确保拿到完整数据)
            byte[] jsonBytes = new byte[dataLength];
            int totalRead = 0;
            while (totalRead < dataLength)
            {
                int currentRead = clientSocketStream.Read(jsonBytes, totalRead, dataLength - totalRead);
                if (currentRead == 0)
                {
                    Console.WriteLine("连接已中断");
                    break;
                }
                totalRead += currentRead;
            }

            // 步骤3:反序列化消息
            using (MemoryStream ms = new MemoryStream(jsonBytes))
            {
                Message m = Message.ReadJson(ms);
                Console.WriteLine("收到消息:" + m);
            }

            // Socket保持打开,可继续等待下一条消息
            Console.WriteLine("Socket保持连接,等待后续数据...");
            // 可添加后续接收或发送逻辑
            // ...

            Console.ReadLine();
            // 不再立即关闭Socket,按需处理
            // clientSocket.Shutdown(SocketShutdown.Both);
            // clientSocket.Close();
            return;
        }
        catch (SocketException se)
        {
            Console.WriteLine(se.Message);
        }
    }
}

关键说明

  • 帧分割:通过先发送长度的方式,给TCP流添加了消息边界,接收端能准确判断单条消息的结束位置。
  • 循环读取:NetworkStream.Read()可能只返回部分数据,必须循环读取直到获取完整的长度和消息内容。
  • Flush():发送数据后调用Flush(),确保缓冲区中的数据立即发送到网络,避免延迟。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 19:15:01