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

Spring WebFlux中如何读取请求体?网关拦截场景实践需求

我之前在做Spring WebFlux网关的时候也踩过这个坑!直接订阅请求体的话,后面的过滤器链或者业务逻辑就拿不到请求体了——因为Reactive里的请求体是个只能消费一次的Flux/Mono。下面给你一套可行的解决方案,亲测有效:

核心思路

要同时实现读取请求体打日志和拦截返回错误,关键是先把请求体缓存起来,避免消费一次后就无法复用。具体步骤:

  1. 缓存原始请求体,生成一个可重复读取的请求对象
  2. 读取缓存的请求体完成日志记录
  3. 根据业务规则判断是否需要拦截返回错误
  4. 若无需拦截,则将缓存后的请求传给后续过滤器链

完整代码实现

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 08:45:54