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:内置空闲检测
✅ 自定义心跳:灵活控制
✅ 断线重连:提高可用性
✅ 连接管理:统一管理所有连接
关键要点
- 心跳间隔建议 30-60 秒
- 心跳超时建议 3 倍心跳间隔
- 使用 IdleStateHandler 简化开发
- 实现断线重连提高可用性
- 定期清理无效连接
上一章:第11章:零拷贝与高性能优化
下一章:第13章:流量整形与限流