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

WebClient结合Http Interface记录流量遇空指针异常求助

问题:WebClient结合Http Interface日志过滤器抛出NullPointerException

我在使用WebClient结合Http Interface时,需要将出入站流量记录到数据库。编写日志过滤器后部分服务抛出NullPointerException,尝试用subscribe替代block也无效,已尝试多种方法仍未解决,Spring Boot版本为3.1.1。

我的过滤器代码

public ExchangeFilterFunction logFilter() {

    return (request, next) ->
        next.exchange(request)
            .publishOn(Schedulers.boundedElastic())
            .doOnNext(
                clientResponse -> {
                  clientResponse
                      .bodyToMono(String.class)
                      .doOnNext(
                          clientResponseAsString -> {
                            try {
                              // save request and response to database
                            } catch (Exception e) {
                              log.error(
                                  "Something went wrong in saveOutboundLog method", e.getCause());
                            }
                          })
                      .block();
                });
  }

我的客户端Bean代码

@Bean
public BaseApiClient client(HttpClient httpClient) {
    final var webClient =
        WebClient.builder()
            .codecs(
                clientCodecConfigurer ->
                    clientCodecConfigurer.defaultCodecs().maxInMemorySize(32 * 1024 * 1024))
            .baseUrl(clientProperties.getBaseUrl())
            .clientConnector(new ReactorClientHttpConnector(httpClient))
            .filter(logFilter())
            .build();
    final var factory =
        HttpServiceProxyFactory.builder(WebClientAdapter.forClient(webClient))
            .blockTimeout(Duration.ofSeconds(60))
            .build();
    return factory.createClient(MyClient.class);
  }

问题原因及解决方案

核心问题

  1. block()破坏响应式流:在doOnNext里调用block()会阻塞线程,且当响应体已被下游Http Interface自动解析时,clientResponse.bodyToMono(String.class)会返回空Mono,调用block()直接抛出NullPointerException。
  2. 响应体只能消费一次:WebClient的ClientResponse响应体是一次性资源,下游Http Interface解析后,过滤器再读取就会失败。
  3. 错误处理覆盖不全:原代码仅捕获了数据库保存异常,未处理响应体读取失败的情况。

修复后的过滤器代码

import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.web.reactive.function.client.ClientRequest;
import org.springframework.web.reactive.function.client.ClientResponse;
import org.springframework.web.reactive.function.client.ExchangeFilterFunction;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;

import java.nio.charset.StandardCharsets;

public ExchangeFilterFunction logFilter() {
    return (request, next) -> {
        // 缓存请求体,支持多次读取
        ClientRequest cachedRequest = ClientRequest.from(request)
                .body(BodyInserters.fromDataBuffers(request.bodyToFlux(DataBuffer.class).cache()))
                .build();

        return next.exchange(cachedRequest)
                .flatMap(clientResponse -> {
                    // 缓存响应体,避免下游和过滤器竞争资源
                    Flux<DataBuffer> responseBody = clientResponse.bodyToFlux(DataBuffer.class).cache();

                    // 读取响应体字符串,空响应返回空串
                    return responseBody.collect(DataBufferUtils::join)
                            .map(dataBuffer -> {
                                byte[] bytes = new byte[dataBuffer.readableByteCount()];
                                dataBuffer.read(bytes);
                                DataBufferUtils.release(dataBuffer);
                                return new String(bytes, StandardCharsets.UTF_8);
                            })
                            .defaultIfEmpty("")
                            .flatMap(responseBodyStr -> {
                                // 异步保存日志,失败不影响主流程
                                return saveLog(cachedRequest, clientResponse, responseBodyStr)
                                        .onErrorResume(e -> {
                                            log.error("保存出入站日志失败", e);
                                            return Mono.empty();
                                        })
                                        .then(Mono.just(clientResponse.mutate()
                                                .body(responseBody)
                                                .build()));
                            });
                })
                .publishOn(Schedulers.boundedElastic()); // 若数据库操作是阻塞型,用此调度器隔离
    };
}

// 替换为实际的日志保存逻辑,推荐用响应式客户端(如R2DBC)
private Mono<Void> saveLog(ClientRequest request, ClientResponse response, String responseBody) {
    // 示例:若用JDBC阻塞操作,需确保在boundedElastic调度下执行
    return Mono.fromRunnable(() -> {
        // 执行数据库插入操作
        // save request method, url, headers and response status, body to database
    });
}

关键优化点

  • 缓存请求/响应体:用cache()确保请求和响应体可多次消费,避免下游与过滤器冲突。
  • 全程响应式操作:移除block(),用flatMap、defaultIfEmpty等操作符保持异步非阻塞特性。
  • 隔离阻塞操作:如果数据库操作是阻塞型(如JDBC),通过publishOn(Schedulers.boundedElastic())隔离,避免阻塞Netty线程池。
  • 重建响应对象:缓存响应体后重新构建ClientResponse,保证下游Http Interface能正常解析响应。
  • 容错处理:用onErrorResume捕获日志保存异常,不影响主业务流程。

内容的提问来源于stack exchange,提问作者Yusuf Erdoğan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 05:33:20