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

如何用Spring Webflux和Reactor Netty检测客户端连接断开?

问题分析与解决

当服务器与客户端断开连接时无法捕获异常,大概率是以下几个原因导致的:

1. handle操作可能吞掉了异常

你的handleJsonResponse方法如果在处理JSON时出现异常,但没有调用sink.error(Throwable)将异常传递下去,就会导致异常被静默吞掉,下游的doOnError自然无法触发。

检查handleJsonResponse的实现,确保任何异常都通过sink.error()抛出:

private <T> void handleJsonResponse(String jsonString, Class<T> entityClass, SynchronousSink<T> sink) {
    try {
        // 你的JSON解析逻辑
        T entity = objectMapper.readValue(jsonString, entityClass);
        sink.next(entity);
    } catch (Exception e) {
        // 必须将异常传递给sink
        sink.error(new DataAccessException("解析JSON失败", e) {});
    }
}

2. retrieve()无法捕获底层IO异常

retrieve()方法只会处理HTTP响应状态码为4xx/5xx的情况,像连接断开、网络超时这类底层IO异常,需要通过exchangeToFlux()手动处理,它能覆盖更全面的错误场景:

修改select方法的请求逻辑:

@Override
public <T> Flux<T> select(Query query, Class<T> entityClass) throws DataAccessException {
    return client.post()
            .uri("/query-stream")
            .contentType(MediaType.APPLICATION_JSON)
            .body(BodyInserters.fromValue("{\"sql\": \"" + sqlQuery + "\"}"))
            .exchangeToFlux(clientResponse -> {
                // 处理HTTP错误状态码
                if (clientResponse.statusCode().isError()) {
                    return clientResponse.createException().flatMapMany(Flux::error);
                }
                // 正常响应转Flux
                return clientResponse.bodyToFlux(String.class);
            })
            .skip(1)
            .<T>handle((jsonString, sink) -> handleJsonResponse(jsonString, entityClass, sink))
            .doOnError(this::handleError)
            .doOnCancel(() -> System.out.println("连接已取消"))
            .doFinally(signalType -> {
                if (signalType == SignalType.CANCEL) {
                    System.out.println("收到取消信号");
                } else if (signalType == SignalType.ERROR) {
                    System.out.println("收到错误信号");
                } else if (signalType == SignalType.COMPLETE) {
                    System.out.println("流正常完成");
                }
            });
}

3. 未配置WebClient超时机制

如果连接断开后没有触发超时,WebClient可能会一直等待,不会抛出异常。给WebClient添加连接超时和读取超时配置:

@Bean
public WebClient webClient() {
    HttpClient httpClient = HttpClient.create()
            .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000) // 连接超时5秒
            .responseTimeout(Duration.ofSeconds(10)); // 读取响应超时10秒

    return WebClient.builder()
            .clientConnector(new ReactorClientHttpConnector(httpClient))
            .build();
}

4. 确认订阅端的错误处理逻辑

你的订阅代码已经包含了错误处理器,但要确保上游的异常能正常传递到这里。如果前面的步骤都正确,连接断开时应该会触发错误处理器里的逻辑,打印异常并释放latch。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 17:54:49