Java线程如何同时等待客户端消息与其他线程消息?
解决方案
要实现线程同时等待客户端消息和其他线程转发的消息,核心是给每个客户端线程分配一个阻塞队列接收转发消息,同时拆分IO读取与消息发送逻辑,避免单阻塞操作占用线程。
关键设计思路
- 每个客户端线程(
ClientHandler)维护一个BlockingQueue,用于存储其他线程需要发送给当前客户端的消息。 - 用子线程专门读取客户端输入流,避免阻塞主线程处理队列消息的逻辑。
- 服务端维护线程安全集合管理所有客户端线程,用于消息广播。
修改后的完整代码
服务端管理类 ChatServer
import java.io.IOException; import java.net.ServerSocket; import java.net.Socket; import java.util.concurrent.CopyOnWriteArrayList; public class ChatServer { private static final int PORT = 12345; // 线程安全集合,管理所有客户端处理器 private static CopyOnWriteArrayList<ClientHandler> clientHandlers = new CopyOnWriteArrayList<>(); public static void main(String[] args) throws IOException { ServerSocket serverSocket = new ServerSocket(PORT); System.out.println("服务端已启动,监听端口:" + PORT); while (true) { Socket clientSocket = serverSocket.accept(); ClientHandler handler = new ClientHandler(clientSocket); clientHandlers.add(handler); handler.start(); } } // 广播消息给所有客户端,排除发送者 public static void broadcastMessage(byte[] message, ClientHandler sender) { for (ClientHandler handler : clientHandlers) { if (handler != sender) { handler.sendToClient(message); } } } // 移除断开连接的客户端处理器 public static void removeClientHandler(ClientHandler handler) { clientHandlers.remove(handler); } }
客户端线程类 ClientHandler
import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; import java.net.Socket; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; public class ClientHandler extends Thread { private Socket clientSocket; private InputStream inputStream; private OutputStream outputStream; // 存储待发送给当前客户端的消息 private BlockingQueue<byte[]> outgoingMessages; public ClientHandler(Socket socket) throws IOException { this.clientSocket = socket; this.inputStream = socket.getInputStream(); this.outputStream = socket.getOutputStream(); // 初始化阻塞队列,容量可根据需求调整 this.outgoingMessages = new ArrayBlockingQueue<>(100); } @Override public void run() { // 启动子线程读取客户端发来的消息 new Thread(this::readClientMessages).start(); // 主线程阻塞等待队列中的消息,发送给客户端 try { while (!clientSocket.isClosed()) { byte[] message = outgoingMessages.take(); // 阻塞等待消息 outputStream.write(message); outputStream.flush(); } } catch (IOException | InterruptedException e) { System.out.println("客户端断开连接:" + clientSocket.getInetAddress()); } finally { closeResources(); ChatServer.removeClientHandler(this); } } // 读取客户端消息并触发广播 private void readClientMessages() { byte[] buffer = new byte[1024]; int bytesRead; try { while ((bytesRead = inputStream.read(buffer)) != -1) { // 复制有效消息(避免buffer残留数据) byte[] message = new byte[bytesRead]; System.arraycopy(buffer, 0, message, 0, bytesRead); // 广播给其他客户端 ChatServer.broadcastMessage(message, this); } } catch (IOException e) { // 客户端断开连接,退出循环 } finally { closeResources(); ChatServer.removeClientHandler(this); } } // 其他线程调用此方法,将消息放入当前客户端的发送队列 public void sendToClient(byte[] message) { try { outgoingMessages.put(message); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } // 关闭所有IO资源 private void closeResources() { try { inputStream.close(); outputStream.close(); clientSocket.close(); } catch (IOException e) { e.printStackTrace(); } } }
方案说明
- 双等待逻辑实现:
- 客户端消息:子线程调用
inputStream.read()阻塞等待客户端发送消息,读取后触发广播。 - 转发消息:主线程调用
outgoingMessages.take()阻塞等待其他线程放入的消息,取出后发送给客户端。
- 客户端消息:子线程调用
- 线程安全保障:
- 用
CopyOnWriteArrayList管理客户端线程,避免并发遍历和修改的异常。 BlockingQueue本身是线程安全的,无需额外同步。
- 用
- 资源清理:无论客户端正常断开还是异常退出,都会自动关闭IO资源并从服务端集合中移除。
内容的提问来源于stack exchange,提问作者Raul Salgado
相关产品推荐
相关产品推荐

