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); }
问题原因及解决方案
核心问题
block()破坏响应式流:在doOnNext里调用block()会阻塞线程,且当响应体已被下游Http Interface自动解析时,clientResponse.bodyToMono(String.class)会返回空Mono,调用block()直接抛出NullPointerException。- 响应体只能消费一次:WebClient的
ClientResponse响应体是一次性资源,下游Http Interface解析后,过滤器再读取就会失败。 - 错误处理覆盖不全:原代码仅捕获了数据库保存异常,未处理响应体读取失败的情况。
修复后的过滤器代码
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
相关产品推荐
相关产品推荐

