Spring WebFlux中如何读取请求体?网关拦截场景实践需求
我之前在做Spring WebFlux网关的时候也踩过这个坑!直接订阅请求体的话,后面的过滤器链或者业务逻辑就拿不到请求体了——因为Reactive里的请求体是个只能消费一次的Flux/Mono。下面给你一套可行的解决方案,亲测有效:
核心思路
要同时实现读取请求体打日志和拦截返回错误,关键是先把请求体缓存起来,避免消费一次后就无法复用。具体步骤:
- 缓存原始请求体,生成一个可重复读取的请求对象
- 读取缓存的请求体完成日志记录
- 根据业务规则判断是否需要拦截返回错误
- 若无需拦截,则将缓存后的请求传给后续过滤器链
完整代码实现
import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferUtils; import org.springframework.http.HttpStatus; import org.springframework.http.MediaType; import org.springframework.http.server.reactive.ServerHttpRequest; import org.springframework.http.server.reactive.ServerHttpRequestDecorator; import org.springframework.http.server.reactive.ServerHttpResponse; import org.springframework.web.server.ServerWebExchange; import org.springframework.web.server.WebFilter; import org.springframework.web.server.WebFilterChain; import reactor.core.publisher.Mono; import java.nio.charset.StandardCharsets; public class RequestLoggingFilter implements WebFilter { private final boolean enabled; private final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(RequestLoggingFilter.class); public RequestLoggingFilter(boolean enabled) { this.enabled = enabled; } @Override public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) { if (!enabled) { logger.debug("Gateway filter is disabled, routing request directly"); return chain.filter(exchange); } logger.debug("Gateway is enabled, processing request: {}", exchange.getRequest().getPath()); ServerHttpRequest originalRequest = exchange.getRequest(); // 合并请求体为单个DataBuffer,方便读取 return DataBufferUtils.join(originalRequest.getBody()) .flatMap(dataBuffer -> { // 复制缓存请求体,保留引用计数避免被提前释放 DataBuffer cachedBuffer = dataBuffer.retainedDuplicate(); // 读取请求体字节数组 byte[] bodyBytes = new byte[dataBuffer.readableByteCount()]; dataBuffer.read(bodyBytes); // 释放原始DataBuffer,避免内存泄漏 DataBufferUtils.release(dataBuffer); // 打印请求体日志(可根据需求格式化) String requestBody = new String(bodyBytes, StandardCharsets.UTF_8); logger.info("Request path: {}, Request body: {}", originalRequest.getPath(), requestBody); // 这里替换成你的业务拦截判断逻辑 boolean shouldBlockRequest = checkIfRequestNeedsBlock(originalRequest, requestBody); if (shouldBlockRequest) { // 返回错误响应给客户端 ServerHttpResponse response = exchange.getResponse(); response.setStatusCode(HttpStatus.BAD_REQUEST); response.getHeaders().setContentType(MediaType.APPLICATION_JSON); String errorResponse = "{\"code\":400,\"message\":\"Invalid request content\"}"; DataBuffer errorBuffer = response.bufferFactory().wrap(errorResponse.getBytes(StandardCharsets.UTF_8)); return response.writeWith(Mono.just(errorBuffer)); } // 装饰原始请求,让后续链能读取缓存的请求体 ServerHttpRequest decoratedRequest = new ServerHttpRequestDecorator(originalRequest) { @Override public Flux<DataBuffer> getBody() { return Flux.just(cachedBuffer); } }; // 将装饰后的请求传给后续过滤器链 return chain.filter(exchange.mutate().request(decoratedRequest).build()); }) // 处理读取请求体时的异常 .onErrorResume(throwable -> { logger.error("Failed to read request body", throwable); ServerHttpResponse response = exchange.getResponse(); response.setStatusCode(HttpStatus.INTERNAL_SERVER_ERROR); return response.setComplete(); }); } /** * 自定义拦截判断逻辑示例 * 可根据请求路径、请求体内容、请求头等判断是否需要拦截 */ private boolean checkIfRequestNeedsBlock(ServerHttpRequest request, String requestBody) { // 示例:如果请求体包含"invalid"字符串则拦截 return requestBody.contains("invalid"); } }
关键注意事项
- 缓存请求体的内存问题:如果你的网关会处理大请求体(比如上传文件),直接缓存整个请求体可能导致内存溢出。这种情况下可以考虑流式记录日志(只记录前N个字符),或者跳过大请求体的日志记录。
- DataBuffer的引用计数:一定要用
retainedDuplicate()复制缓存,否则原始DataBuffer被释放后,后续链读取会报错。使用完不需要的DataBuffer要调用DataBufferUtils.release()释放。 - 字符编码:示例中用了UTF-8,实际可以从请求头的
Content-Type中获取编码,避免乱码。 - 错误响应的完整性:返回错误时要确保设置正确的状态码和响应头,并且完成响应写入,不要继续调用
chain.filter()。
内容的提问来源于stack exchange,提问作者user9349304
相关产品推荐
相关产品推荐

