Spring Gateway全局过滤器引发Connection prematurely closed问题求助
解决Spring Cloud Gateway POST请求Body处理引发的连接提前关闭问题
你的问题出在请求Body被消费后未重新注入请求:WebFlux中请求的Body是一次性的Flux流,一旦通过doOnNext订阅消费,后续过滤器链或下游服务就无法再获取请求体,直接导致请求中断,抛出PrematureCloseException: Connection prematurely closed BEFORE response错误。
下面是两种可行的修复方案:
方案一:手动缓存并重新包装请求体
@Component public class GlobalPayloadFilter implements GlobalFilter { @Override public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) { HttpMethod method = exchange.getRequest().getMethod(); if (method != HttpMethod.POST) { return chain.filter(exchange); } // 合并所有DataBuffer为一个,缓存请求体 return DataBufferUtils.join(exchange.getRequest().getBody()) .flatMap(dataBuffer -> { // 保留数据引用,防止被提前释放 DataBufferUtils.retain(dataBuffer); // 创建可重复订阅的缓存Body流 Flux<DataBuffer> cachedBody = Flux.defer(() -> Flux.just(dataBuffer.slice(0, dataBuffer.readableByteCount()))); // 构建新请求,替换Body为缓存流 ServerHttpRequest newRequest = exchange.getRequest().mutate().body(cachedBody).build(); ServerWebExchange newExchange = exchange.mutate().request(newRequest).build(); // 处理缓存的Body数据 byte[] content = new byte[dataBuffer.readableByteCount()]; dataBuffer.read(content); DataBufferUtils.release(dataBuffer); // 释放资源,避免内存泄漏 /* 在这里处理content,比如解析、打印等 */ return chain.filter(newExchange); }); } }
方案二:使用Gateway工具类简化缓存逻辑
@Component public class GlobalPayloadFilter implements GlobalFilter { @Override public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) { HttpMethod method = exchange.getRequest().getMethod(); if (method != HttpMethod.POST) { return chain.filter(exchange); } // 利用Gateway内置工具类缓存请求体 return ServerWebExchangeUtils.cacheRequestBody(exchange, cachedRequest -> { // 收集缓存的Body数据 return cachedRequest.getBody().collectList().flatMap(dataBuffers -> { // 合并所有DataBuffer为字节数组 byte[] content = dataBuffers.stream() .map(buffer -> { byte[] bytes = new byte[buffer.readableByteCount()]; buffer.read(bytes); DataBufferUtils.release(buffer); return bytes; }) .reduce(new byte[0], (a, b) -> { byte[] merged = new byte[a.length + b.length]; System.arraycopy(a, 0, merged, 0, a.length); System.arraycopy(b, 0, merged, a.length, b.length); return merged; }); /* 在这里处理content数据 */ // 继续执行过滤器链 return chain.filter(exchange.mutate().request(cachedRequest).build()); }); }); } }
注意事项
- 必须调用
DataBufferUtils.release()释放处理后的DataBuffer,避免内存泄漏。 - 如果处理大文件上传场景,建议使用流式处理代替缓存整个Body,防止内存溢出。
内容的提问来源于stack exchange,提问作者André
相关产品推荐
相关产品推荐

