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

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();
    }
}
问题排查与解决方案

问题根源

  1. 客户端重复创建流导致协议混乱:UserClient的每个请求方法都重复创建ObjectOutputStream和ObjectInputStream。ObjectOutputStream初始化时会写入流头,多次创建会导致服务端的ObjectInputStream解析流结构时出错,触发EOFException。
  2. 服务端嵌套循环耗尽输入流: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 22:43:11