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

Socket简易聊天服务器:如何仅在输入流有数据时读取并广播?

问题描述

我正尝试实现一个极其简易的聊天客户端与服务器,目的是熟悉Socket的使用,而非打造一个合理的聊天系统,但仍希望符合基本常识。本次问题聚焦于服务器代码。

假设已建立连接(accept()方法返回),服务器需为该客户端提供服务(按规范使用独立线程):

public class SimpleChatServer {
    // possibly something else
    private static final ExecutorService executorService = Executors.newCachedThreadPool();
    private static final List<OutputStream> outputStreamList = new ArrayList<>();

    public static void main(String[] args) {
        start();
    }

    @SneakyThrows
    private static void start() {
        try (var serverSocket = new ServerSocket(5000)) {
            while (true) {
                Socket socket = serverSocket.accept();
                executorService.submit(() -> serviceClient(socket));
            }
        }
    }

我需要一个合理的serviceClient()实现。第一步是注册该Socket的输出流,因为服务器需要读取传入的字符串并向所有连接的客户端“广播”:

private static void serviceClient(Socket socket) {
        registerOutputStream(socket);
        // some other code I haven't written yet
    }
private static void registerOutputStream(Socket socket) throws IOException {
        // so that I "foreach" it once I receive a message to broadcast
        outputStreamList.add(socket.getOutputStream());
    }

现在我陷入了困境:服务器如何知晓客户端是否写入了数据,从而可以从Socket的输入流读取?我需要实现的逻辑是:“如果输入流有数据,就读取并广播给客户端;如果没有,则持续等待直到有数据”。这听起来像是一个带有频繁条件检查的无限循环(效率低下)。或者让客户端通知服务器“我已写入数据,你可以读取”,这可以用并发同步机制实现,但让客户端与服务器共享同步对象的设计很奇怪(即使是如此简易的应用)。合理的解决方案是什么?

我认为“你是否了解异步IO”并非充分答案。如果异步IO确实是我应考虑的方案,能否提供一些细节(可能是简单的代码片段)?我对该包并不熟悉。

解决方案

一、同步阻塞IO(入门首选)

Socket的输入流本身就是阻塞式的,不需要自己做轮询检查。当调用read()相关方法时,如果输入流没有数据,线程会自动进入等待状态,直到有数据到来或者连接关闭——这完全符合你需要的“没数据就等,有数据就读取”的逻辑,且不会做无用的循环检查,效率很高。

另外需要补充几个基础细节:

  • outputStreamList是多线程共享资源,必须做线程安全保护,否则会触发并发修改异常
  • 客户端断开连接时,要从列表中移除对应的输出流,避免广播时报错
  • 读取聊天消息建议用BufferedReader按行读取,符合消息传输的常规场景

修改后的完整serviceClient()及配套方法示例:

// 改用线程安全的CopyOnWriteArrayList,避免并发修改问题
private static final List<OutputStream> outputStreamList = new CopyOnWriteArrayList<>();

private static void serviceClient(Socket socket) {
    try (InputStream inputStream = socket.getInputStream();
         BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream, StandardCharsets.UTF_8))) {
        // 注册当前客户端的输出流
        registerOutputStream(socket);
        String message;
        // readLine()是阻塞方法:无数据时线程等待,客户端断开则返回null
        while ((message = reader.readLine()) != null) {
            // 广播消息给所有在线客户端
            broadcastMessage(message);
        }
    } catch (IOException e) {
        // 客户端异常断开时会进入此处
        e.printStackTrace();
    } finally {
        // 客户端断开后清理资源
        removeOutputStream(socket);
        try {
            socket.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

private static void registerOutputStream(Socket socket) throws IOException {
    outputStreamList.add(socket.getOutputStream());
}

private static void removeOutputStream(Socket socket) {
    outputStreamList.removeIf(os -> {
        try {
            // 通过输出流关联的Socket判断是否为当前客户端的流
            return ((SocketOutputStream) os).getSocket() == socket;
        } catch (Exception e) {
            return false;
        }
    });
}

private static void broadcastMessage(String message) throws IOException {
    byte[] messageBytes = (message + System.lineSeparator()).getBytes(StandardCharsets.UTF_8);
    for (OutputStream os : outputStreamList) {
        os.write(messageBytes);
        os.flush(); // 确保消息立即发送到客户端
    }
}

二、异步IO(Java NIO 极简示例)

如果想尝试异步IO,可以用Java NIO的Selector机制,它能让单个线程管理多个连接的IO事件,无需为每个连接单独开线程。以下是极简实现示例:

public class NioChatServer {
    public static void main(String[] args) throws IOException {
        Selector selector = Selector.open();
        ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
        serverSocketChannel.bind(new InetSocketAddress(5000));
        serverSocketChannel.configureBlocking(false);
        // 注册"接受新连接"的事件
        serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);

        while (true) {
            // 阻塞等待IO事件触发(新连接、有数据可读等)
            selector.select();
            Set<SelectionKey> selectedKeys = selector.selectedKeys();
            Iterator<SelectionKey> iterator = selectedKeys.iterator();

            while (iterator.hasNext()) {
                SelectionKey key = iterator.next();
                iterator.remove();

                if (key.isAcceptable()) {
                    // 处理新客户端连接
                    ServerSocketChannel server = (ServerSocketChannel) key.channel();
                    SocketChannel clientChannel = server.accept();
                    clientChannel.configureBlocking(false);
                    // 注册"读取数据"的事件
                    clientChannel.register(selector, SelectionKey.OP_READ);
                    System.out.println("新客户端连接");
                } else if (key.isReadable()) {
                    // 处理客户端发送的消息
                    SocketChannel clientChannel = (SocketChannel) key.channel();
                    ByteBuffer buffer = ByteBuffer.allocate(1024);
                    int bytesRead = clientChannel.read(buffer);
                    if (bytesRead == -1) {
                        // 客户端断开连接,清理资源
                        clientChannel.close();
                        key.cancel();
                        continue;
                    }
                    buffer.flip();
                    String message = StandardCharsets.UTF_8.decode(buffer).toString().trim();
                    // 广播消息给所有在线客户端
                    broadcastMessage(selector, message);
                }
            }
        }
    }

    private static void broadcastMessage(Selector selector, String message) throws IOException {
        byte[] messageBytes = (message + System.lineSeparator()).getBytes(StandardCharsets.UTF_8);
        ByteBuffer buffer = ByteBuffer.wrap(messageBytes);
        for (SelectionKey key : selector.keys()) {
            Channel channel = key.channel();
            if (channel instanceof SocketChannel && key.isValid()) {
                SocketChannel clientChannel = (SocketChannel) channel;
                clientChannel.write(buffer);
                buffer.rewind(); // 重置缓冲区,用于给下一个客户端发送消息
            }
        }
    }
}

这个NIO版本中,selector.select()会阻塞直到有IO事件发生,无需轮询,也不需要为每个连接开线程,适合连接数较多的场景,但对于入门学习Socket来说,同步阻塞IO更容易理解和实现。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 20:05:57