如何实现Spring WebFlux中byte[]与Flux<DataBuffer>双向转换
Spring WebFlux 过滤器提前读取请求体实现方案
你要的byte[]和FluxDataBufferUtils实现,以下是可直接落地的生产可用代码,完全满足「业务逻辑消费前提前读取请求体、处理后不影响下游正常读取」的需求。
完整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
相关产品推荐
相关产品推荐

