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

如何实现Spring WebFlux中byte[]与Flux<DataBuffer>双向转换

Spring WebFlux 过滤器提前读取请求体实现方案

你要的byte[]和Flux双向转换逻辑可以直接通过DataBufferUtils实现,以下是可直接落地的生产可用代码,完全满足「业务逻辑消费前提前读取请求体、处理后不影响下游正常读取」的需求。

完整Filter实现代码

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

import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.nio.channels.Channels;

@Component
public class PreReadRequestBodyFilter implements WebFilter {

    @Override
    public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
        // 跳过无请求体的请求,可按需扩展过滤规则,比如跳过GET请求、文件上传请求
        long contentLength = exchange.getRequest().getHeaders().getContentLength();
        if (contentLength == 0) {
            return chain.filter(exchange);
        }

        // 1. 聚合请求体Flux<DataBuffer>为byte[]
        return DataBufferUtils.join(exchange.getRequest().getBody())
                .flatMap(originalDataBuffer -> {
                    byte[] bodyBytes;
                    try (ByteArrayOutputStream outputStream = new ByteArrayOutputStream()) {
                        // 读取DataBuffer内容到字节数组,用只读视图避免修改原始缓冲区
                        Channels.newChannel(outputStream)
                                .write(originalDataBuffer.asByteBuffer().asReadOnlyBuffer());
                        bodyBytes = outputStream.toByteArray();
                    } catch (IOException e) {
                        return Mono.error(new RuntimeException("读取请求体失败", e));
                    } finally {
                        // 必须手动释放join拿到的原始DataBuffer,否则会造成堆外内存泄漏
                        DataBufferUtils.release(originalDataBuffer);
                    }

                    // 2. 执行自定义逻辑,比如签名校验、日志打印、请求参数校验等
                    handleRequestBody(bodyBytes, exchange);

                    // 3. 基于byte[]重建请求体,通过装饰器回写到exchange
                    ServerHttpRequest decoratedRequest = new ServerHttpRequestDecorator(exchange.getRequest()) {
                        @Override
                        public Flux<DataBuffer> getBody() {
                            // 使用当前exchange绑定的缓冲区工厂创建新的DataBuffer,兼容不同web容器
                            DataBuffer newBuffer = exchange.getResponse().bufferFactory()
                                    .allocateBuffer(bodyBytes.length)
                                    .write(bodyBytes);
                            return Flux.just(newBuffer);
                        }
                    };

                    return chain.filter(exchange.mutate().request(decoratedRequest).build());
                });
    }

    private void handleRequestBody(byte[] bodyBytes, ServerWebExchange exchange) {
        // 替换为你的自定义业务逻辑即可,示例为打印请求体
        String charset = exchange.getRequest().getHeaders().getContentType().getCharset().name();
        String bodyContent = new String(bodyBytes, charset);
        System.out.printf("路径[%s]提前读取到请求体:%s%n", 
                exchange.getRequest().getPath().value(), bodyContent);
    }
}

关键注意事项

  • 内存泄漏风险:DataBufferUtils.join()返回的DataBuffer是堆外内存分配的,必须在finally块中手动调用DataBufferUtils.release()释放,否则会出现堆外内存泄漏,线上服务会直接出现内存溢出
  • 容器兼容性:重建DataBuffer时必须使用exchange.getResponse().bufferFactory()获取缓冲区工厂实例,不要手动构造DataBuffer,否则在Netty、Jetty、Tomcat等不同容器下可能出现类型不兼容问题
  • 大请求适配:DataBufferUtils.join()会把整个请求体加载到内存,如果你的服务存在大体积请求体的场景,需要在配置文件中调整spring.codec.max-in-memory-size参数,设置合理的内存阈值,避免大请求打垮服务
  • 性能优化:可以根据业务规则过滤不需要提前读取的请求,比如静态资源请求、GET请求、文件上传请求,减少不必要的内存拷贝开销

实现逻辑说明

常见的请求体日志打印方案是在原有请求体Flux上添加doOnNext、map等旁路算子,只有当下游业务逻辑开始消费请求体时才会触发读取,属于消费时拦截,无法满足提前读取的需求。
上述方案在过滤器执行阶段就主动聚合完整请求体,在下游业务逻辑执行前就拿到全量字节数组完成自定义处理,之后重新构造一个全新的请求体流传递给下游,对后续业务逻辑完全透明,不会出现请求体被消费后下游读不到内容的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:48:14