diff --git a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyWriteResponseFilter.java b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyWriteResponseFilter.java index 3dae9d95..278c428b 100644 --- a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyWriteResponseFilter.java +++ b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyWriteResponseFilter.java @@ -23,6 +23,7 @@ import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import reactor.core.publisher.SignalType; import reactor.netty.Connection; import org.springframework.core.Ordered; @@ -98,8 +99,12 @@ public class NettyWriteResponseFilter implements GlobalFilter, Ordered { return (isStreamingMediaType(contentType) ? response.writeAndFlushWith(body.map(Flux::just)) : response.writeWith(body)); - })).doOnCancel(() -> cleanup(exchange)) - .doOnError(throwable -> cleanup(exchange)); + })) + .doFinally(signalType -> { + if (signalType == SignalType.CANCEL || signalType == SignalType.ON_ERROR) { + cleanup(exchange); + } + }); // @formatter:on } @@ -121,7 +126,7 @@ public class NettyWriteResponseFilter implements GlobalFilter, Ordered { private void cleanup(ServerWebExchange exchange) { Connection connection = exchange.getAttribute(CLIENT_RESPONSE_CONN_ATTR); - if (connection != null && connection.channel().isActive()) { + if (connection != null) { connection.dispose(); } }