ServerSocket与客户端通信抛出StreamCorruptedException问题咨询
解决Java NIO服务端与客户端通信的StreamCorruptedException错误
我在SelectionKey处于可写状态时通过ServerSocket向客户端发送消息,服务端的发送逻辑看似正确,但客户端读取时抛出以下错误:
java.io.StreamCorruptedException: invalid stream header: 48656C6C
at java.io.ObjectInputStream.readStreamHeader(ObjectInputStream.java:866)
at java.io.ObjectInputStream.(ObjectInputStream.java:358)
at MyTcpClient.main(MyTcpClient.java:22)
请问客户端该如何正确读取消息?相关代码如下:
MyAsyncProcessor.java
import java.io.IOException; import java.net.InetAddress; import java.net.InetSocketAddress; import java.nio.ByteBuffer; import java.nio.CharBuffer; import java.nio.channels.*; import java.nio.charset.CharsetEncoder; import java.nio.charset.StandardCharsets; import java.util.HashMap; import java.util.Iterator; import java.util.Set; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class MyAsyncProcessor { HashMap<Integer, MyTask> hashMap = new HashMap<>(); ExecutorService pool; public MyAsyncProcessor() { } public static void main(String[] args) throws IOException { new MyAsyncProcessor().process(); } public void process() throws IOException { pool = Executors.newFixedThreadPool(2); InetAddress host = InetAddress.getByName("localhost"); Selector selector = Selector.open(); ServerSocketChannel serverSocketChannel = ServerSocketChannel.open(); serverSocketChannel.configureBlocking(false); serverSocketChannel.bind(new InetSocketAddress(host, 9876)); serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT); SelectionKey key; System.out.println("AsyncProcessor ready"); while (true) { if (selector.select() > 0) { Set<SelectionKey> selectedKeys = selector.selectedKeys(); Iterator<SelectionKey> i = selectedKeys.iterator(); while (i.hasNext()) { key = i.next(); if (!key.isValid()) { key.cancel(); continue; } if (key.isAcceptable()) { SocketChannel socketChannel = serverSocketChannel.accept(); socketChannel.configureBlocking(false); System.out.println("Channel hashCode: " + socketChannel.hashCode()); MyTask task = new MyTask(selector, socketChannel); hashMap.put(socketChannel.hashCode(), task); pool.execute(task); socketChannel.register(selector, SelectionKey.OP_READ); } if (key.isReadable()) { SocketChannel socketChannel = (SocketChannel) key.channel(); MyTask task = hashMap.get(socketChannel.hashCode()); ByteBuffer byteBuffer = ByteBuffer.allocate(1024); try { socketChannel.read(byteBuffer); String result = new String(byteBuffer.array()).trim(); String[] words = result.split(" "); task.timeToRead = Integer.parseInt(words[words.length - 2]) * 1000; task.timeToWrite = Integer.parseInt(words[words.length - 1]) * 1000; System.out.println(Thread.currentThread().getName() + " reads for " + task.timeToRead + " mills"); try { Thread.sleep(task.timeToRead); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println(Thread.currentThread().getName() + " ended reading"); } catch (Exception e) { e.printStackTrace(); return; } socketChannel.register(selector, SelectionKey.OP_WRITE); } if (key.isWritable()) { SocketChannel socketChannel = (SocketChannel) key.channel(); MyTask task = hashMap.get(socketChannel.hashCode()); task.readyToWrite(); hashMap.remove(socketChannel.hashCode()); System.out.println(Thread.currentThread().getName() + " writes for " + task.timeToWrite + " mills"); try { Thread.sleep(task.timeToWrite); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println(Thread.currentThread().getName() + " ended writing"); CharsetEncoder enc = StandardCharsets.US_ASCII.newEncoder(); String response = "Hello!\n"; socketChannel.write(enc.encode(CharBuffer.wrap(response))); key.cancel(); } i.remove(); } } } } }
MyTcpClient.java
import java.io.*; import java.net.Socket; import java.util.Random; public class MyTcpClient { public static void main(String[] args) { Random rand = new Random(); int secondsToRead = rand.nextInt(5); int secondsToWrite = secondsToRead + 1; String message = "Seconds for the task to be read and written: " + secondsToRead + " " + secondsToWrite; System.out.println(message); Socket socket; try { socket = new Socket("127.0.0.1", 9876); PrintWriter printWriter = new PrintWriter(socket.getOutputStream(), true); printWriter.println(message); System.out.println("Sending message"); ObjectInputStream ois = new ObjectInputStream(socket.getInputStream()); String response = (String) ois.readObject(); System.out.println("Response: " + response); ois.close(); } catch (IOException e) { System.out.println("Error in Socket"); e.printStackTrace(); System.exit(-1); } catch (ClassNotFoundException e) { throw new RuntimeException(e); } } }
问题根源
服务端发送的是普通ASCII文本数据,但客户端误用ObjectInputStream读取——这个类是专门用来读取Java序列化对象的,它会先检查固定的序列化流头(0xACED0005),而服务端发送的"Hello!"的前四个字节是48 65 6C 6C(对应ASCII的H、e、l、l),完全不符合序列化流头格式,因此抛出错误。
修复方案
客户端需要改用与服务端匹配的文本读取方式,比如用BufferedReader读取整行文本:
修改MyTcpClient.java的读取逻辑:
// 替换原ObjectInputStream相关代码 BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream(), StandardCharsets.US_ASCII)); String response = reader.readLine(); System.out.println("Response: " + response); reader.close();
额外优化提示
- 服务端发送的响应包含换行符
\n,客户端用readLine()可以刚好读取完整的响应内容,协议匹配度高。 - 如果后续需要传输Java对象,服务端要改用
ObjectOutputStream写入序列化对象,客户端再用ObjectInputStream读取,保持两端协议一致。 - 服务端的NIO读取逻辑存在数据截断风险:
socketChannel.read(byteBuffer)可能只读取了部分客户端消息,建议循环读取直到获取完整数据(比如读取到换行符),避免解析错误。
内容的提问来源于stack exchange,提问作者Luca
相关产品推荐
相关产品推荐

