You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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();
    }
}

关键优化点说明

  1. 移除自定义Handler:将计时逻辑通过doOnNext和doOnError实现,避免引入额外的Netty Handler,简化代码结构
  2. 直接在onErrorResume处理异常:无需手动触发pipeline异常事件,直接在响应式流中完成异常判断,逻辑更清晰
  3. 明确非特定异常处理:主动调用connection.dispose()关闭连接,避免资源泄漏,替代原方案中“不做操作”的模糊处理
  4. 线程安全的计时:用ThreadLocal替代AtomicLong,避免同一连接下多请求并发时的计时混乱

内容的提问来源于stack exchange,提问作者bscott

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 05:13:11