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
相关产品推荐
相关产品推荐

