Netty TCP Server如何在抛出自定义异常时保持连接不关闭?
问题
用Project Reactor Netty开发TCP长连接服务器,接收客户端字节数组请求,处理后返回字节数组响应,仅在致命错误时关闭连接。遇到场景:处理请求抛出特定异常时,需保持连接且不返回数据,但当前抛出这类异常时Netty会自动关闭连接。尝试自定义Handler的exceptionCaught(...)方法但从未触发,设置ChannelOption.AUTO_CLOSE=false也无效(该选项仅适用于写入异常,此场景无回写)。
临时解决方案
以下是确保exceptionCaught()触发、处理异常并保持连接的临时实现:
服务器初始化代码
DisposableServer someTcpServer = tcpServer .host("12.123.456.789") .port(12345) .wiretap(true) .doOnBind(server -> log.info("Starting listener...")) .doOnBound(server -> log.info("Listener started on host: {}, port: {}", server.host(), server.port())) .option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT) .childOption(ChannelOption.AUTO_CLOSE,false) .childOption(ChannelOption.SO_KEEPALIVE, true) .doOnConnection(connection -> { InetSocketAddress socketAddress = (InetSocketAddress) connection.channel().remoteAddress(); log.info("Client has connected. Host: {}, Port: {}", socketAddress.getAddress().getHostAddress(), socketAddress.getPort()); }) .doOnChannelInit((observer, channel, remoteAddress) -> channel.pipeline() .addFirst(new TcpServerHandler()) ) .handle((inbound, outbound) -> inbound .receive() .asByteArray() .flatMap(req -> processRequest(req)) .flatMap(rsp -> outbound.sendByteArray(Flux.just(rsp))) // 捕获processRequest及后续抛出的异常 .onErrorResume(throwable -> { // 手动触发pipeline的异常处理 inbound.withConnection(connection -> connection.channel().pipeline().fireExceptionCaught(throwable)); return Mono.empty(); }) ).bindNow(); someTcpServer.onDispose().block();
自定义TcpServerHandler代码
@Slf4j public class TcpServerHandler extends ChannelDuplexHandler { private final AtomicLong startTime = new AtomicLong(0L); private final AtomicLong endTime = new AtomicLong(0L); @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { if (cause instanceof BusinessLogicException1 || cause instanceof BusinessLogicException2) { endTime.set(System.nanoTime()); log.info("Took {} ms to process", Duration.ofNanos(endTime.get() - startTime.get()).toMillis()); // 特定异常:保持连接,无需额外操作 ctx.fireChannelActive(); } else { // 非特定异常:原调用super导致方法触发两次,故移除后不做操作 } } @Override public void channelRead(ChannelHandlerContext ctx, Object msg) { startTime.set(System.nanoTime()); ctx.fireChannelRead(msg); } }
方案验证与优化建议
当前方案的合理性
你的临时方案是可行的:
- 通过
onErrorResume捕获处理流程中的异常,手动触发pipeline的异常事件,确保自定义Handler能感知到异常 - 针对特定异常不执行关闭操作,Netty不会主动断开长连接
- 移除
super.exceptionCaught(...)避免重复触发是正确的,因为ChannelDuplexHandler的父类ChannelInboundHandlerAdapter的exceptionCaught会默认将异常传递给下一个Handler,可能导致重复处理
更优实现方式
可以简化流程,无需手动触发pipeline异常,直接在onErrorResume中完成异常判断、计时和处理,同时保留必要的连接维护逻辑:
优化后的服务器初始化代码
DisposableServer someTcpServer = tcpServer .host("12.123.456.789") .port(12345) .wiretap(true) .doOnBind(server -> log.info("Starting listener...")) .doOnBound(server -> log.info("Listener started on host: {}, port: {}", server.host(), server.port())) .option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT) .childOption(ChannelOption.SO_KEEPALIVE, true) .doOnConnection(connection -> { InetSocketAddress socketAddress = (InetSocketAddress) connection.channel().remoteAddress(); log.info("Client has connected. Host: {}, Port: {}", socketAddress.getAddress().getHostAddress(), socketAddress.getPort()); }) .handle((inbound, outbound) -> inbound .receive() .asByteArray() .doOnNext(req -> { // 将计时逻辑移到此处,替代Handler的channelRead TcpServerMetrics.startRequest(); }) .flatMap(req -> processRequest(req)) .flatMap(rsp -> outbound.sendByteArray(Flux.just(rsp))) .doOnError(throwable -> { // 异常时计算耗时 long durationMs = TcpServerMetrics.endRequest(); log.info("Took {} ms to process", durationMs); }) .onErrorResume(throwable -> { if (throwable instanceof BusinessLogicException1 || throwable instanceof BusinessLogicException2) { // 特定异常:返回空Mono,不发送响应,保持连接 return Mono.empty(); } else { // 非特定异常:关闭连接并传递异常 return inbound.withConnection(connection -> { connection.dispose(); return Mono.error(throwable); }); } }) ).bindNow(); someTcpServer.onDispose().block();
辅助计时工具类(替代Handler的计时逻辑)
public class TcpServerMetrics { private static final ThreadLocal<Long> REQUEST_START_TIME = new ThreadLocal<>(); public static void startRequest() { REQUEST_START_TIME.set(System.nanoTime()); } public static long endRequest() { long end = System.nanoTime(); long start = REQUEST_START_TIME.get(); REQUEST_START_TIME.remove(); return Duration.ofNanos(end - start).toMillis(); } }
关键优化点说明
- 移除自定义Handler:将计时逻辑通过
doOnNext和doOnError实现,避免引入额外的Netty Handler,简化代码结构 - 直接在
onErrorResume处理异常:无需手动触发pipeline异常事件,直接在响应式流中完成异常判断,逻辑更清晰 - 明确非特定异常处理:主动调用
connection.dispose()关闭连接,避免资源泄漏,替代原方案中“不做操作”的模糊处理 - 线程安全的计时:用
ThreadLocal替代AtomicLong,避免同一连接下多请求并发时的计时混乱
内容的提问来源于stack exchange,提问作者bscott
相关产品推荐
相关产品推荐

