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

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é

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 09:25:35