Netty教程 / 第 120 节

第12章:心跳检测与空闲连接管理

本章导读

心跳检测是保持长连接稳定的重要机制。本章将讲解为什么需要心跳、IdleStateHandler 的使用、自定义心跳实现和断线重连机制。


12.1 为什么需要心跳机制

12.1.1 问题场景

问题1:假死连接

  • 网络异常导致连接断开
  • 但应用层不知道,仍然保持连接
  • 浪费服务器资源

问题2:NAT超时

  • NAT设备会清除长时间空闲的连接
  • 导致连接断开

问题3:防火墙

  • 防火墙可能关闭长时间空闲的连接

12.1.2 解决方案

心跳机制

  • 定期发送心跳包
  • 检测连接是否存活
  • 及时清理死连接

12.2 IdleStateHandler 空闲检测

12.2.1 基本用法

public class HeartbeatServer {
    public static void main(String[] args) throws Exception {
        EventLoopGroup bossGroup = new NioEventLoopGroup(1);
        EventLoopGroup workerGroup = new NioEventLoopGroup();
        
        try {
            ServerBootstrap bootstrap = new ServerBootstrap();
            bootstrap.group(bossGroup, workerGroup)
                    .channel(NioServerSocketChannel.class)
                    .childHandler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        protected void initChannel(SocketChannel ch) {
                            ChannelPipeline pipeline = ch.pipeline();
                            
                            // 空闲检测Handler
                            pipeline.addLast(new IdleStateHandler(
                                5,   // readerIdleTime: 读空闲时间(秒)
                                10,  // writerIdleTime: 写空闲时间(秒)
                                15   // allIdleTime: 读写空闲时间(秒)
                            ));
                            
                            // 业务Handler
                            pipeline.addLast(new HeartbeatServerHandler());
                        }
                    });
            
            ChannelFuture future = bootstrap.bind(8080).sync();
            System.out.println("心跳服务器启动");
            future.channel().closeFuture().sync();
        } finally {
            bossGroup.shutdownGracefully();
            workerGroup.shutdownGracefully();
        }
    }
}

class HeartbeatServerHandler extends SimpleChannelInboundHandler<String> {
    
    @Override
    protected void channelRead0(ChannelHandlerContext ctx, String msg) {
        System.out.println("收到消息: " + msg);
        
        if ("PING".equals(msg)) {
            ctx.writeAndFlush("PONG");
        }
    }
    
    @Override
    public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
        if (evt instanceof IdleStateEvent) {
            IdleStateEvent event = (IdleStateEvent) evt;
            
            String eventType = null;
            switch (event.state()) {
                case READER_IDLE:
                    eventType = "读空闲";
                    // 5秒没有读取到数据,关闭连接
                    System.out.println("读空闲,关闭连接");
                    ctx.close();
                    break;
                case WRITER_IDLE:
                    eventType = "写空闲";
                    break;
                case ALL_IDLE:
                    eventType = "读写空闲";
                    break;
            }
            
            System.out.println(ctx.channel().remoteAddress() + " " + eventType);
        }
    }
}

12.2.2 客户端心跳

public class HeartbeatClient {
    public static void main(String[] args) throws Exception {
        EventLoopGroup group = new NioEventLoopGroup();
        
        try {
            Bootstrap bootstrap = new Bootstrap();
            bootstrap.group(group)
                    .channel(NioSocketChannel.class)
                    .handler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        protected void initChannel(SocketChannel ch) {
                            ChannelPipeline pipeline = ch.pipeline();
                            
                            pipeline.addLast(new StringDecoder());
                            pipeline.addLast(new StringEncoder());
                            
                            // 空闲检测
                            pipeline.addLast(new IdleStateHandler(0, 3, 0));
                            
                            // 业务Handler
                            pipeline.addLast(new HeartbeatClientHandler());
                        }
                    });
            
            ChannelFuture future = bootstrap.connect("localhost", 8080).sync();
            System.out.println("连接服务器成功");
            future.channel().closeFuture().sync();
        } finally {
            group.shutdownGracefully();
        }
    }
}

class HeartbeatClientHandler extends SimpleChannelInboundHandler<String> {
    
    @Override
    protected void channelRead0(ChannelHandlerContext ctx, String msg) {
        System.out.println("收到响应: " + msg);
    }
    
    @Override
    public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
        if (evt instanceof IdleStateEvent) {
            IdleStateEvent event = (IdleStateEvent) evt;
            
            if (event.state() == IdleState.WRITER_IDLE) {
                // 3秒没有写数据,发送心跳
                System.out.println("发送心跳: PING");
                ctx.writeAndFlush("PING");
            }
        }
    }
}

12.3 自定义心跳实现

public class CustomHeartbeatHandler extends ChannelInboundHandlerAdapter {
    
    private static final int HEARTBEAT_INTERVAL = 30;  // 心跳间隔(秒)
    private static final int HEARTBEAT_TIMEOUT = 90;   // 心跳超时(秒)
    
    private ScheduledFuture<?> heartbeatFuture;
    private long lastHeartbeatTime;
    
    @Override
    public void channelActive(ChannelHandlerContext ctx) {
        lastHeartbeatTime = System.currentTimeMillis();
        
        // 启动心跳任务
        heartbeatFuture = ctx.executor().scheduleAtFixedRate(() -> {
            long currentTime = System.currentTimeMillis();
            long timeSinceLastHeartbeat = currentTime - lastHeartbeatTime;
            
            if (timeSinceLastHeartbeat > HEARTBEAT_TIMEOUT * 1000) {
                // 超时,关闭连接
                System.out.println("心跳超时,关闭连接");
                ctx.close();
            } else {
                // 发送心跳
                System.out.println("发送心跳");
                ctx.writeAndFlush("HEARTBEAT");
            }
        }, HEARTBEAT_INTERVAL, HEARTBEAT_INTERVAL, TimeUnit.SECONDS);
        
        ctx.fireChannelActive();
    }
    
    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) {
        // 更新最后心跳时间
        lastHeartbeatTime = System.currentTimeMillis();
        
        if ("HEARTBEAT".equals(msg)) {
            // 收到心跳,回复
            ctx.writeAndFlush("HEARTBEAT_ACK");
        } else {
            ctx.fireChannelRead(msg);
        }
    }
    
    @Override
    public void channelInactive(ChannelHandlerContext ctx) {
        // 取消心跳任务
        if (heartbeatFuture != null) {
            heartbeatFuture.cancel(true);
        }
        
        ctx.fireChannelInactive();
    }
}

12.4 断线重连机制

public class ReconnectClient {
    
    private final String host;
    private final int port;
    private final Bootstrap bootstrap;
    private Channel channel;
    
    public ReconnectClient(String host, int port) {
        this.host = host;
        this.port = port;
        
        EventLoopGroup group = new NioEventLoopGroup();
        bootstrap = new Bootstrap();
        bootstrap.group(group)
                .channel(NioSocketChannel.class)
                .handler(new ChannelInitializer<SocketChannel>() {
                    @Override
                    protected void initChannel(SocketChannel ch) {
                        ch.pipeline().addLast(new ReconnectHandler(ReconnectClient.this));
                    }
                });
    }
    
    public void connect() {
        try {
            ChannelFuture future = bootstrap.connect(host, port).sync();
            channel = future.channel();
            System.out.println("连接成功");
        } catch (Exception e) {
            System.err.println("连接失败: " + e.getMessage());
            scheduleReconnect();
        }
    }
    
    public void scheduleReconnect() {
        bootstrap.config().group().schedule(() -> {
            System.out.println("尝试重连...");
            connect();
        }, 5, TimeUnit.SECONDS);
    }
    
    public static void main(String[] args) {
        ReconnectClient client = new ReconnectClient("localhost", 8080);
        client.connect();
    }
}

class ReconnectHandler extends ChannelInboundHandlerAdapter {
    
    private final ReconnectClient client;
    
    public ReconnectHandler(ReconnectClient client) {
        this.client = client;
    }
    
    @Override
    public void channelInactive(ChannelHandlerContext ctx) {
        System.out.println("连接断开,准备重连");
        client.scheduleReconnect();
    }
    
    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
        System.err.println("发生异常: " + cause.getMessage());
        ctx.close();
    }
}

12.5 连接管理最佳实践

public class ConnectionManager {
    
    private final Map<String, Channel> connections = new ConcurrentHashMap<>();
    private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
    
    public void addConnection(String id, Channel channel) {
        connections.put(id, channel);
        
        // 监听连接关闭
        channel.closeFuture().addListener(future -> {
            connections.remove(id);
            System.out.println("连接关闭: " + id);
        });
    }
    
    public void removeConnection(String id) {
        Channel channel = connections.remove(id);
        if (channel != null && channel.isActive()) {
            channel.close();
        }
    }
    
    public void broadcast(String message) {
        connections.values().forEach(channel -> {
            if (channel.isActive()) {
                channel.writeAndFlush(message);
            }
        });
    }
    
    public void startHealthCheck() {
        scheduler.scheduleAtFixedRate(() -> {
            connections.entrySet().removeIf(entry -> {
                Channel channel = entry.getValue();
                if (!channel.isActive()) {
                    System.out.println("清理无效连接: " + entry.getKey());
                    return true;
                }
                return false;
            });
        }, 60, 60, TimeUnit.SECONDS);
    }
}

12.6 本章小结

心跳机制:保持长连接稳定
IdleStateHandler:内置空闲检测
自定义心跳:灵活控制
断线重连:提高可用性
连接管理:统一管理所有连接

关键要点

  1. 心跳间隔建议 30-60 秒
  2. 心跳超时建议 3 倍心跳间隔
  3. 使用 IdleStateHandler 简化开发
  4. 实现断线重连提高可用性
  5. 定期清理无效连接

上一章第11章:零拷贝与高性能优化
下一章第13章:流量整形与限流