IO、NIO、Netty实战
·
目标
客户端和服务端互相通信,本文主要是实战练习,照着敲,然后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();
}
});
}
}
更多推荐




所有评论(0)