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

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();
        }
    }
}

方案说明

  1. 双等待逻辑实现:
    • 客户端消息:子线程调用inputStream.read()阻塞等待客户端发送消息,读取后触发广播。
    • 转发消息:主线程调用outgoingMessages.take()阻塞等待其他线程放入的消息,取出后发送给客户端。
  2. 线程安全保障:
    • 用CopyOnWriteArrayList管理客户端线程,避免并发遍历和修改的异常。
    • BlockingQueue本身是线程安全的,无需额外同步。
  3. 资源清理:无论客户端正常断开还是异常退出,都会自动关闭IO资源并从服务端集合中移除。

内容的提问来源于stack exchange,提问作者Raul Salgado

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 08:05:09