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

