目标

客户端和服务端互相通信,本文主要是实战练习,照着敲,然后debug看为什么就行

前置理解

模型 核心类 特点 简述
BIO ServerSocket / Socket 一个连接一个线程,accept()read() 都会阻塞 简单但连接多了线程爆炸
NIO Selector / ServerSocketChannel / SocketChannel / ByteBuffer 一个 Selector 管理多个 Channel,非阻塞轮询 复杂但一个线程扛几百连接
Netty ServerBootstrap / EventLoopGroup / Pipeline / Handler 封装了 NIO 的复杂性,用责任链处理数据 高性能框架,写业务逻辑就行

练习递进关系:BIO 理解通信本质 → NIO 理解多路复用 → Netty 理解工程化封装。

三种服务端可以混搭客户端:底层都是 TCP 协议,BIO 客户端可以连 NIO 服务端,反之亦然。

IO

服务端

public class IOServer {  
    private static final ThreadPoolExecutor  clientPool = new ThreadPoolExecutor (  
            2, //最小线程数  
            3, //最大线程数  
            60L, //线程空闲时间  
            TimeUnit.SECONDS, // 时间单位  
            new ArrayBlockingQueue<>(1), // 任务队列容量  
            new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略  
    );  
    public static void main(String[] args) throws IOException {  
        //IO服务端  
        ServerSocket serverSocket = new ServerSocket(8888);  
        while (true) {  
            // 添加线程池状态监控  
            System.out.println("当前线程池状态 - 活跃线程: " + clientPool.getActiveCount() +  
                    ", 队列大小: " + clientPool.getQueue().size());  
  
            Socket clientSocket = serverSocket.accept();  
            System.out.println("客户端连接: " + clientSocket.getInetAddress());  
//            new Thread(()->handleClient(clientSocket)).start();  
            try {  
                // 提交任务  
                clientPool.execute(() -> handleClient(clientSocket));  
                System.out.println("✅ 任务已提交");  
  
            } catch (Exception e) {  
                System.err.println("❌ 提交失败");  
                clientSocket.close();  
            }  
        }  
  
  
    }  
  
    private static void handleClient(Socket clientSocket){  
        try {  
            clientSocket.setSoTimeout(60000);  
        } catch (SocketException e) {  
            return;  
        }  
        //inpustStream获取输入流  
        //outputStream获取输出流  
        //InputStreamReader将输入流转换为字符流  
        //BufferedReader将字符流转换为缓冲流  
  
        try(InputStream inputStream = clientSocket.getInputStream();  
            BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream));  
           ){  
            String line;  
            while ((line = reader.readLine()) != null) {  
                System.out.printf(Thread.currentThread().getName() + "收到客户端消息:"+line+"\n");  
            }  
            System.out.println("客户端已断开");  
        }catch (SocketTimeoutException e) {  
            System.out.println("客户端空闲超时,主动踢下线");  
        } catch (IOException e) {  
            e.printStackTrace();  
        }finally {  
            try {  
                clientSocket.close(); // 只关客户端,不关服务端!  
                System.out.println("资源已清理,线程退出");  
            } catch (IOException e) {  
                e.printStackTrace();  
            }  
        }  
    }  
}

客户端

public class IOClient {  
    public static void main(String[] args) throws IOException {  
        //IO客户端  
        Socket socket = new Socket("127.0.0.1",8888);  
        OutputStream os = socket.getOutputStream();  
        Scanner sc = new Scanner(System.in);  
  
        while (sc.hasNext()){  
            String msg = sc.nextLine();  
            if ("exit".equalsIgnoreCase(msg)) {  
                socket.shutdownOutput(); // 关闭输出流,通知服务端读完了  
                break;  
            }  
            os.write((msg + "\n").getBytes());  
            os.flush();  
        }  
        os.close();  
        socket.close();  
    }  
}

NIO

服务端

public class NIOServer {  
    public static final ThreadPoolExecutor tp = new ThreadPoolExecutor(  
            2,10,60L, TimeUnit.SECONDS,  
            new ArrayBlockingQueue<>(1),  
            new ThreadPoolExecutor.CallerRunsPolicy()  
    );  
  
    private NIOServer start () throws IOException {  
        //配置nio管理员  
        Selector selector = Selector.open();  
        ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();  
        serverSocketChannel.bind(new InetSocketAddress(8888));  
        serverSocketChannel.configureBlocking(false);   //取消阻塞  
  
        // 注册  
        serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);  
        while (true){  
            selector.select();  
            Set<SelectionKey> selectionKeys = selector.selectedKeys();  
            //获取迭代器  
            Iterator<SelectionKey> iterator = selectionKeys.iterator();  
            while (iterator.hasNext()){  
                SelectionKey key = iterator.next();  
                iterator.remove();  
                // 处理不同事件  
                if (key.isAcceptable()) {  
                    handleAccept(key); // 新客户端连接  
                } else if (key.isReadable()) {  
                    handleRead(key);   // 客户端发消息  
                }  
            }  
        }  
    }  
    public static void main(String[] args) throws IOException {  
       new NIOServer().start();  
  
    }  
  
    private static void handleAccept(SelectionKey key) throws IOException {  
        ServerSocketChannel server = (ServerSocketChannel) key.channel();  
        SocketChannel channel = server.accept();  
        channel.configureBlocking(false);  
        channel.register(key.selector(),SelectionKey.OP_READ, ByteBuffer.allocate(1024));  
        System.out.println("\n新客户端连接:" + channel.getRemoteAddress());  
  
  
    }  
  
    private static void handleRead(SelectionKey key){  
        SocketChannel client = (SocketChannel) key.channel();  
        ByteBuffer buffer = (ByteBuffer) key.attachment();  
  
        try {  
            int len = client.read(buffer); // 非阻塞读取  
            if (len == -1) {  
                // 客户端断开  
                System.out.println("客户端断开:" + client.getRemoteAddress());  
                client.close();  
                key.cancel();  
                return;  
            }  
  
            // 切换读模式  
            buffer.flip();  
            String msg = new String(buffer.array(), 0, len);  
            buffer.clear();  
            System.out.println("收到消息:" + msg);  
            tp.execute(() -> {  
                System.out.println(Thread.currentThread().getName() + " 执行业务 → " + msg);  
            });  
  
        } catch (IOException e) {  
            // 强制断开  
            try {  
                client.close();  
                key.cancel();  
            } catch (IOException ex) {  
  
                ex.printStackTrace();  
            }  
        }  
    }  
}

客户端

public class NIOClient {
    public static void main(String[] args) throws IOException {
        // 1. 打开客户端通道
        SocketChannel client = SocketChannel.open(new InetSocketAddress("127.0.0.1", 8888));
        client.configureBlocking(false);  // 非阻塞模式
        
        Scanner sc = new Scanner(System.in);
        ByteBuffer buffer = ByteBuffer.allocate(1024);

        while (sc.hasNext()) {
            String msg = sc.nextLine();
            if ("exit".equalsIgnoreCase(msg)) {
                break;
            }
            // 写数据到服务端
            buffer.clear();
            buffer.put((msg + "\n").getBytes());
            buffer.flip();
            client.write(buffer);
        }
        client.close();
    }
}

Netty

客户端

public class NettyClient {  
    private Channel channel;  
    public final EventLoopGroup group = new NioEventLoopGroup();  
    public static void main(String[] args) throws InterruptedException {  
        NettyClient client = new NettyClient();  
        new Thread(()->{  
            try {  
                client.start("127.0.0.1",8888);  
            } catch (InterruptedException e) {  
                throw new RuntimeException(e);  
            }  
        }).start();  
        Scanner sc = new Scanner(System.in);  
        while (sc.hasNext()){  
            String msg = sc.nextLine();  
            if("exit".equalsIgnoreCase(msg)) {  
                client.sendMsg("拜拜");  
                client.close();  
                break;  
            }  
            client.sendMsg(msg);  
        }  
    }  
  
    public void start(String host,Integer port) throws InterruptedException {  
        Bootstrap bootstrap = new Bootstrap();  
        bootstrap  
                .group(group)  
                .channel(NioSocketChannel.class)  
                .handler(new ChannelInitializer<SocketChannel>() {  
  
                    @Override  
                    protected void initChannel(SocketChannel ch) throws Exception {  
                        ch.pipeline()  
//                                .addLast(new LoggingHandler(LogLevel.INFO))  
                                // 拆包参数和服务端一模一样  
                                .addLast(new LengthFieldBasedFrameDecoder(1024 * 1024, 0, 4, 0, 4))  
                                .addLast(new LengthFieldPrepender(4))  
                                .addLast(new StringDecoder())  
                                .addLast(new StringEncoder())  
                                // 客户端心跳  
                                .addLast(new IdleStateHandler(30, 60, 120))  
                                // 客户端自定义业务处理器  
                                .addLast(new MyClientHandler());  
  
                    }  
                });  
        ChannelFuture future = bootstrap.connect(host, port).sync();  
        this.channel = future.channel();  
        future.channel().closeFuture().sync();  
    }  
    public void sendMsg(String msg) {  
        if (channel != null && channel.isActive()) {  
            channel.writeAndFlush(msg);  
        } else {  
            System.out.println("连接未成功,无法发送");  
        }  
    }  
    public void close() {  
        if (channel != null) channel.close();  
        group.shutdownGracefully();  
    }  
}
方法配置
public class MyClientHandler extends SimpleChannelInboundHandler<String> {  
    @Override  
    protected void channelRead0(ChannelHandlerContext ctx, String msg) throws Exception {  
        System.out.println("客户端收到服务端消息:" + msg);  
    }  
    // 连接成功建立  
    @Override  
    public void channelActive(ChannelHandlerContext ctx) {  
        System.out.println("客户端成功连接服务端");  
        // 连接成功主动发一条消息  
        ctx.writeAndFlush("客户端上线啦");  
    }  
  
    // 连接断开  
    @Override  
    public void channelInactive(ChannelHandlerContext ctx) {  
        System.out.println("客户端与服务端连接断开");  
    }  
  
    // 心跳事件处理  
    @Override  
    public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {  
        // 判断是不是空闲事件  
        if (evt instanceof IdleStateEvent) {  
  
            IdleStateEvent event = (IdleStateEvent) evt;  
            // 如果是【写空闲】= 30秒没发消息  
            if (event.state() == IdleState.WRITER_IDLE) {  
                // 自动发一个心跳包!  
                ctx.writeAndFlush("heartbeat");  
            }  
        }  
    }  
}

服务端

public class NettyServer {  
    private static final EventLoopGroup boss = new NioEventLoopGroup(1);  
    private static final EventLoopGroup worker  = new NioEventLoopGroup();  
    public void start(int port) throws InterruptedException {  
        try {  
            ServerBootstrap bootstrap = new ServerBootstrap();  
            bootstrap  
                    .group(boss, worker)  
                    .channel(NioServerSocketChannel.class)  
                    .childHandler(new ChannelInitializer<SocketChannel>() {  
                        @Override  
                        protected void initChannel(SocketChannel ch) throws Exception {  
                            //输入输出  
                            ChannelPipeline pipeline = ch.pipeline();  
//                            pipeline.addLast(new LoggingHandler(LogLevel.INFO)); //日志打印  
                            pipeline.addLast(new LengthFieldBasedFrameDecoder(1024 * 1024, 0, 4, 0, 4));  
                            pipeline.addLast(new LengthFieldPrepender(4));  
                            pipeline.addLast(new StringDecoder());  
                            pipeline.addLast(new StringEncoder());  
                            pipeline.addLast(new IdleStateHandler(60, 30, 120));  
                            pipeline.addLast(new MyServerHandler());  
                        }  
                    });  
            //绑定端口 + 同步阻塞  
            ChannelFuture channelFuture = bootstrap.bind(port).sync();  
            //等待服务关闭通道  
            channelFuture.channel().closeFuture().sync();  
        }catch(Exception e){  
            e.printStackTrace();  
        }finally {  
            boss.shutdownGracefully();  
            worker.shutdownGracefully();  
        }  
    }  
    public static void main(String[] args) throws InterruptedException {  
        // 启动服务端  
        new NettyServer().start(8888);  
    }  
}
方法配置
public class MyServerHandler extends SimpleChannelInboundHandler<String> {  
    //全局用户集合,concurrentHashMap  
    public static final Set<String> ONLINE_CLIENTS = ConcurrentHashMap.newKeySet();  
    private static final ThreadPoolExecutor tp = new ThreadPoolExecutor(5, 10, 60L,  
            TimeUnit.MILLISECONDS, new ArrayBlockingQueue<>(1),new ThreadPoolExecutor.AbortPolicy());  
  
    @Override  
    public void channelActive(ChannelHandlerContext ctx) {  
        SocketAddress socketAddress = ctx.channel().remoteAddress();  
        tp.execute(() -> {  
            ONLINE_CLIENTS.add(socketAddress.toString());  
            System.out.println("有新客户端连接:" + ctx.channel().remoteAddress());  
            System.out.println("当前人数:"+ONLINE_CLIENTS.size());  
            ctx.writeAndFlush("欢迎连接!");  
            ctx.writeAndFlush("当前人数:"+ONLINE_CLIENTS.size());  
        });  
    }  
    @Override  
    public void channelReadComplete(ChannelHandlerContext ctx) {  
        ctx.flush();  
    }  
  
    @Override  
    public void channelInactive(ChannelHandlerContext ctx) {  
        String clientAddr = ctx.channel().remoteAddress().toString();  
        tp.execute(() -> {  
            System.out.println("客户端断开:" + clientAddr);  
            ONLINE_CLIENTS.remove(clientAddr);  
            System.out.println("当前在线人数:" + ONLINE_CLIENTS.size());  
  
        });  
    }  
  
    @Override  
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {  
        System.err.println("异常:" + cause.getMessage());  
        ctx.close();  
    }  
    @Override  
    protected void channelRead0(ChannelHandlerContext ctx, String msg) throws Exception {  
        String name = ctx.channel().remoteAddress().toString();  
        tp.execute(() -> {  
            try {  
                System.out.println(name+" 发来一条消息:" + msg);  
                ctx.writeAndFlush("ACK_OK");  
            } catch (Exception e) {  
                e.printStackTrace();  
            }  
        });  
  
    }  
}
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐