Java 21 Socket多连接场景下二次请求触发EOFException问题排查
问题描述
搭建了包含UserClient、BridgeServer(中间层)和InventoryServer的分布式架构,基于Java 21 Socket实现多连接通信,使用ObjectInputStream和ObjectOutputStream传输序列化对象,但出现以下问题:
- 启动BridgeServer;
- 启动UserClient,确认连接已建立;
- 发送
AuthenticationRequest等请求,首次请求处理正常; - 再次发送任意请求时,BridgeServer在
Object request = in.readObject();处抛出EOFException
错误栈信息
java.io.EOFException at java.base/java.io.ObjectInputStream$PeekInputStream.readFully(ObjectInputStream.java:2933) at java.base/java.io.ObjectInputStream$BlockDataInputStream.readShort(ObjectInputStream.java:3428) at java.base/java.io.ObjectInputStream.readStreamHeader(ObjectInputStream.java:985) at java.base/java.io.ObjectInputStream.<init>(ObjectInputStream.java:416) at com.pe.distributed.system.bridge.BridgeServer.run(BridgeSide.java:44) at com.pe.distributed.system.bridge.BridgeMain.main(BridgeMain.java:5)
核心代码
UserClient
import java.io.IOException; import java.io.ObjectInputStream; import java.io.ObjectOutputStream; import java.io.Serializable; import java.net.Socket; import java.util.Map; import java.util.Scanner; import java.util.logging.Level; import java.util.logging.Logger; public record AuthenticationRequest(String username, String password) implements Serializable {} public record InventoryQuery() implements Serializable {} public record ItemOrder(String itemCode, int quantity) implements Serializable {} public class UserClient { private static final Logger logger = Logger.getLogger(UserClient.class.getName()); private static final String EXCEPTION_OCCURRED = "Exception occurred"; private UserClient() {} @SuppressWarnings("java:S2189") public static void run() { Scanner scanner = new Scanner(System.in); while (true) { try (Socket bridgeSocket = new Socket("127.0.0.1", 5550)) { logger.log(Level.INFO, "Connected successfully to bridge server"); while (true) { displayMainMenu(); int choice = scanner.nextInt(); scanner.nextLine(); switch (choice) { case 1: authenticate(scanner, bridgeSocket); break; case 2: checkInventory(bridgeSocket); break; case 3: placeOrder(scanner, bridgeSocket); break; default: logger.log(Level.WARNING, "Invalid choice. Please enter a valid option."); } } } catch (IOException | ClassNotFoundException e) { logger.log(Level.SEVERE, EXCEPTION_OCCURRED, e); } finally { scanner.close(); } } } @SuppressWarnings("java:S106") private static void displayMainMenu() { System.out.println("What would you like to do:"); System.out.println("1. Authenticate"); System.out.println("2. Check Inventory"); System.out.println("3. Place an Order"); System.out.print("Enter your choice: "); } @SuppressWarnings("java:S106") private static void authenticate(Scanner scanner, Socket socket) throws IOException, ClassNotFoundException { Map.Entry<String, String> authCreds = getUserAuthenticationInput(scanner); AuthenticationRequest authenticationRequest = new AuthenticationRequest(authCreds.getKey(), authCreds.getValue()); ObjectOutputStream bridgeIn = new ObjectOutputStream(socket.getOutputStream()); bridgeIn.writeObject(authenticationRequest); bridgeIn.flush(); ObjectInputStream bridgeOut = new ObjectInputStream(socket.getInputStream()); Object response = bridgeOut.readObject(); logger.log(Level.FINE, "Received response from bridge: {0}", response); System.out.println(response); } @SuppressWarnings("java:S106") private static Map.Entry<String, String> getUserAuthenticationInput(Scanner scanner) { System.out.println("Please provide credentials.\n"); System.out.print("Username: "); String username = scanner.nextLine(); System.out.print("Password: "); String password = scanner.nextLine(); return Map.entry(username, password); } @SuppressWarnings("java:S106") private static void checkInventory(Socket socket) throws IOException, ClassNotFoundException { ObjectOutputStream bridgeIn = new ObjectOutputStream(socket.getOutputStream()); bridgeIn.writeObject(new InventoryQuery()); bridgeIn.flush(); ObjectInputStream bridgeOut = new ObjectInputStream(socket.getInputStream()); Object response = bridgeOut.readObject(); logger.log(Level.FINE, "Received response from bridge: {0}", response); System.out.println(response); } private static void placeOrder(Scanner scanner, Socket socket) throws IOException, ClassNotFoundException { Map.Entry<String, Integer> orderDetails = getUserOrderInput(scanner); ItemOrder itemOrder = new ItemOrder(orderDetails.getKey(), orderDetails.getValue()); ObjectOutputStream bridgeIn = new ObjectOutputStream(socket.getOutputStream()); bridgeIn.writeObject(itemOrder); bridgeIn.flush(); ObjectInputStream bridgeOut = new ObjectInputStream(socket.getInputStream()); Object response = bridgeOut.readObject(); logger.log(Level.FINE, "Received response from bridge: {0}", response); } @SuppressWarnings("java:S106") private static Map.Entry<String, Integer> getUserOrderInput(Scanner scanner) { System.out.println("Place your order!\n"); String itemCode; int desiredQuantity; do { System.out.print("Item-code: "); itemCode = scanner.nextLine(); } while (itemCode == null || itemCode.isEmpty()); do { System.out.print("Quantity: "); while (!scanner.hasNextInt()) { System.out.println("Invalid input. Please enter a valid number."); scanner.next(); } desiredQuantity = scanner.nextInt(); scanner.nextLine(); } while (desiredQuantity <= 0); return Map.entry(itemCode, desiredQuantity); } }
BridgeServer
import lombok.Getter; import lombok.Setter; import java.io.IOException; import java.io.ObjectInputStream; import java.io.ObjectOutputStream; import java.io.Serializable; import java.net.ServerSocket; import java.net.Socket; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.logging.Level; import java.util.logging.Logger; @Getter @Setter public class User { String username; String password; public User(String username, String password) { this.username = username; this.password = password; } } public record AuthenticationResponse(String message, Boolean isAuthenticated) implements Serializable {} public record AuthenticationRequest(String username, String password) implements Serializable {} public class BridgeServer { private static final Logger logger = Logger.getLogger(BridgeServer.class.getName()); private static final Map<String, User> users; static { users = new ConcurrentHashMap<>(); users.put("user1", new User("user1", "pass1")); users.put("user2", new User("user2", "pass2")); users.put("user3", new User("user3", "pass3")); } @SuppressWarnings("java:S2189") public static void run() { try (ExecutorService executorService = Executors.newCachedThreadPool(); ServerSocket serverSocket = new ServerSocket(5550)) { logger.info("Bridge server started, waiting for connections..."); while (true) { Socket userSocket = serverSocket.accept(); String clientKey = userSocket.getInetAddress().toString() + ":" + userSocket.getPort(); logger.log(Level.INFO, "User client with key {0} connected.", clientKey); executorService.submit(new BridgeConnectionHandler(userSocket, clientKey)); } } catch (IOException e) { logger.log(Level.SEVERE, "IOException occurred", e); } } @SuppressWarnings("java:S2189") private record BridgeConnectionHandler(Socket userSocket, String clientKey) implements Runnable { @Override public void run() { try (ObjectInputStream in = new ObjectInputStream(userSocket.getInputStream()); ObjectOutputStream out = new ObjectOutputStream(userSocket.getOutputStream())) { boolean isAuthenticated = false; while (true) { Object request = in.readObject(); if (isAuthenticated) { if (!(request instanceof AuthenticationRequest)) { handleUserRequest(in, out); } else { out.writeObject(new AuthenticationResponse("Client already authenticated with user " + users.get(clientKey), true)); out.flush(); } } else if (request instanceof AuthenticationRequest authenticationRequest) { AuthenticationResponse authenticationResponse = handleAuthenticationRequest(authenticationRequest); out.writeObject(authenticationResponse); out.flush(); isAuthenticated = authenticationResponse.isAuthenticated(); } else { out.writeObject(new AuthenticationResponse("Not authenticated", false)); out.flush(); } } } catch (IOException | ClassNotFoundException e) { logger.log(Level.SEVERE, "Exception occurred", e); } } private AuthenticationResponse handleAuthenticationRequest(AuthenticationRequest authenticationRequest) { boolean authenticated = authenticate(authenticationRequest.username(), authenticationRequest.password()); return new AuthenticationResponse(authenticated ? "Authentication successful" : "Authentication failed", authenticated); } private static synchronized boolean authenticate(String username, String password) { User user = users.get(username); return user != null && user.getPassword().equals(password); } private void handleUserRequest(ObjectInputStream bridgeIn, ObjectOutputStream userOut) throws IOException { try (Socket inventorySocket = new Socket("127.0.0.1", 12346); ObjectOutputStream inventoryOut = new ObjectOutputStream(inventorySocket.getOutputStream()); ObjectInputStream inventoryIn = new ObjectInputStream(inventorySocket.getInputStream())) { Object request; while ((request = bridgeIn.readObject()) != null) { logger.log(Level.INFO, "Received request from user: {0}", request); inventoryOut.writeObject(request); inventoryOut.flush(); Object response = inventoryIn.readObject(); logger.log(Level.INFO, "Received response from inventory: {0}", response); userOut.writeObject(response); userOut.flush(); } } catch (ClassNotFoundException e) { logger.log(Level.SEVERE, "Exception occurred", e); } } } public static void main(String[] args) { run(); } }
问题排查与解决方案
问题根源
- 客户端重复创建流导致协议混乱:UserClient的每个请求方法都重复创建
ObjectOutputStream和ObjectInputStream。ObjectOutputStream初始化时会写入流头,多次创建会导致服务端的ObjectInputStream解析流结构时出错,触发EOFException。 - 服务端嵌套循环耗尽输入流:BridgeServer的
handleUserRequest方法内部存在嵌套循环,会持续读取输入流直到EOF,导致第一次认证后的请求被耗尽,外层循环再读取时触发EOFException。
解决方案
1. 客户端:复用Socket的输入输出流
连接建立后只创建一次ObjectOutputStream和ObjectInputStream,所有请求复用这两个流:
public static void run() { Scanner scanner = new Scanner(System.in); while (true) { try (Socket bridgeSocket = new Socket("127.0.0.1", 5550); ObjectOutputStream bridgeOut = new ObjectOutputStream(bridgeSocket.getOutputStream()); ObjectInputStream bridgeIn = new ObjectInputStream(bridgeSocket.getInputStream())) { logger.log(Level.INFO, "Connected successfully to bridge server"); while (true) { displayMainMenu(); int choice = scanner.nextInt(); scanner.nextLine(); switch (choice) { case 1: authenticate(scanner, bridgeOut, bridgeIn); break; case 2: checkInventory(bridgeOut, bridgeIn); break; case 3: placeOrder(scanner, bridgeOut, bridgeIn); break; default: logger.log(Level.WARNING, "Invalid choice. Please enter a valid option."); } } } catch (IOException | ClassNotFoundException e) { logger.log(Level.SEVERE, EXCEPTION_OCCURRED, e); } finally { scanner.close(); } } } // 修正后的authenticate方法 private static void authenticate(Scanner scanner, ObjectOutputStream bridgeOut, ObjectInputStream bridgeIn) throws IOException, ClassNotFoundException { Map.Entry<String, String> authCreds = getUserAuthenticationInput(scanner); AuthenticationRequest authenticationRequest = new AuthenticationRequest(authCreds.getKey(), authCreds.getValue()); bridgeOut.writeObject(authenticationRequest); bridgeOut.flush(); Object response = bridgeIn.readObject(); logger.log(Level.FINE, "Received response from bridge: {0}", response); System.out.println(response); }
同理,checkInventory和placeOrder方法需修改为接收已创建的流对象,不再重复创建。
2. 服务端:移除嵌套循环,处理单个请求
将handleUserRequest改为处理单个请求,外层循环负责持续监听客户端请求:
@Override public void run() { try (ObjectInputStream in = new ObjectInputStream(userSocket.getInputStream()); ObjectOutputStream out = new ObjectOutputStream(userSocket.getOutputStream())) { boolean isAuthenticated = false; while (true) { Object request = in.readObject(); if (isAuthenticated) { if (!(request instanceof AuthenticationRequest)) { handleSingleUserRequest(request, out); } else { out.writeObject(new AuthenticationResponse("Client already authenticated", true)); out.flush(); } } else if (request instanceof AuthenticationRequest authenticationRequest) { AuthenticationResponse authenticationResponse = handleAuthenticationRequest(authenticationRequest); out.writeObject(authenticationResponse); out.flush(); isAuthenticated = authenticationResponse.isAuthenticated(); } else { out.writeObject(new AuthenticationResponse("Not authenticated", false)); out.flush(); } } } catch (IOException | ClassNotFoundException e) { logger.log(Level.SEVERE, "Exception occurred", e); } } private void handleSingleUserRequest(Object request, ObjectOutputStream userOut) throws IOException { try (Socket inventorySocket = new Socket("127.0.0.1", 12346); ObjectOutputStream inventoryOut = new ObjectOutputStream(inventorySocket.getOutputStream()); ObjectInputStream inventoryIn = new ObjectInputStream(inventorySocket.getInputStream())) { logger.log(Level.INFO, "Received request from user: {0}", request); inventoryOut.writeObject(request); inventoryOut.flush(); Object response = inventoryIn.readObject(); logger.log(Level.INFO, "Received response from inventory: {0}", response); userOut.writeObject(response); userOut.flush(); } catch (ClassNotFoundException e) { logger.log(Level.SEVERE, "Exception occurred", e); } }
关键注意事项
- 每个Socket连接只能创建一对
ObjectOutputStream和ObjectInputStream,避免重复写入流头导致解析错误; - 使用try-with-resources自动管理流和Socket的生命周期,防止提前关闭连接;
- 服务端与客户端的请求处理逻辑需一一对应,避免嵌套循环耗尽输入流。
内容的提问来源于stack exchange,提问作者JavaEnthusiast27
相关产品推荐
相关产品推荐

