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

WebFlux WebFilter中解析请求体写入Header及DataBuffer正确关闭方法

WebFlux WebFilter 正确提取请求体并处理 DataBuffer 方案

问题拆解

你当前的代码存在两个核心问题:

  1. DataBuffer 资源泄漏:直接将 DataBuffer 转为 InputStream 后未正确释放,会导致内存泄漏。
  2. 请求体无法复用:原请求体被消费后,后续过滤链或控制器无法再读取请求体;且未实现将提取的 rqUid 放入请求头的需求。

正确实现方案

要解决这些问题,需通过复制请求体保证资源安全与请求体复用,同时修改请求头传递提取的值。以下是完整实现:

import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.http.HttpHeaders;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.http.server.reactive.ServerHttpRequestDecorator;
import org.springframework.web.server.WebFilter;
import org.springframework.web.server.WebFilterChain;
import reactor.core.publisher.Mono;

import java.nio.charset.StandardCharsets;

public class RqUidFilter implements WebFilter {

    private static final String RQUID_HEADER = "X-Rq-Uid";

    @Override
    public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
        ServerHttpRequest originalRequest = exchange.getRequest();

        // 合并请求体的所有 DataBuffer,统一处理资源
        return DataBufferUtils.join(originalRequest.getBody())
                .flatMap(originalBuffer -> {
                    // 创建原始 buffer 的副本,供后续链路读取请求体
                    DataBuffer bufferCopy = originalBuffer.factory().allocateBuffer(originalBuffer.readableByteCount());
                    bufferCopy.write(originalBuffer);
                    // 立即释放原始 buffer,避免资源泄漏
                    DataBufferUtils.release(originalBuffer);

                    // 从副本中解析请求体字符串
                    String requestBody = StandardCharsets.UTF_8.decode(bufferCopy.asByteBuffer()).toString();
                    String rqUid = getRequestRqUid(requestBody);

                    // 包装原始请求,添加目标请求头
                    ServerHttpRequest modifiedRequest = new ServerHttpRequestDecorator(originalRequest) {
                        @Override
                        public HttpHeaders getHttpHeaders() {
                            HttpHeaders updatedHeaders = new HttpHeaders();
                            updatedHeaders.putAll(super.getHttpHeaders());
                            if (rqUid != null) {
                                updatedHeaders.add(RQUID_HEADER, rqUid);
                            }
                            return updatedHeaders;
                        }
                    };

                    // 将副本作为新的请求体传递给后续链路,并在请求结束后释放副本
                    return chain.filter(exchange.mutate().request(modifiedRequest).build())
                            .doFinally(signal -> DataBufferUtils.release(bufferCopy));
                });
    }

    // 根据实际业务逻辑解析请求体中的 rqUid
    private String getRequestRqUid(String requestBody) {
        // 示例:假设请求体是 JSON,使用 Jackson 解析
        // ObjectMapper mapper = new ObjectMapper();
        // try {
        //     return mapper.readTree(requestBody).get("rqUid").asText();
        // } catch (JsonProcessingException e) {
        //     return null;
        // }
        return "parsed_rq_uid_example";
    }
}

关键细节说明

  • DataBuffer 资源管理:
    • 使用 DataBufferUtils.join 合并所有请求体 DataBuffer,避免逐个处理的繁琐。
    • 创建原始 buffer 的副本后立即释放原始资源,最后通过 doFinally 在请求生命周期结束时释放副本,确保无资源泄漏。
  • 请求体复用:通过 ServerHttpRequestDecorator 包装请求,将副本作为新的请求体传递,保证后续组件能正常读取请求体。
  • 请求头修改:在装饰器中重写 getHttpHeaders 方法,添加提取到的 rqUid,让后续链路通过请求头直接获取该值。
  • 上下文传递(可选):如果需要在 Reactor 上下文中也传递 rqUid,可以在 chain.filter 后添加 .subscriberContext(ctx -> ctx.put(RQUID_HEADER, rqUid)),但请求头传递已覆盖大部分场景需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 06:31:06