From 7b1baae2667322dddb7854d09c41aa6854fe9804 Mon Sep 17 00:00:00 2001 From: Maurice Barnum Date: Thu, 2 Mar 2017 14:17:46 -0800 Subject: [PATCH] Close HTTP client connection upon server error or per HTTP spec Besides 500, other errors indicate that the server is having problems, so close the connection to avoid sending requests to the same bad server. DefaultKeepAliveStrategy will look at the request and response headers for other reasons to close, such as "Connection: close" headers. --- .../java/com/yahoo/pulsar/client/impl/HttpClient.java | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/HttpClient.java b/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/HttpClient.java index 4850fb2bb60ca..4f96c4be0e8a3 100644 --- a/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/HttpClient.java +++ b/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/HttpClient.java @@ -26,9 +26,12 @@ import java.util.concurrent.CompletableFuture; import io.netty.channel.EventLoopGroup; +import io.netty.handler.codec.http.HttpRequest; +import io.netty.handler.codec.http.HttpResponse; import io.netty.handler.codec.http.HttpResponseStatus; import io.netty.handler.ssl.SslContext; import org.asynchttpclient.*; +import org.asynchttpclient.channel.DefaultKeepAliveStrategy; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -70,8 +73,12 @@ protected HttpClient(String serviceUrl, Authentication authentication, EventLoop confBuilder.setConnectTimeout(connectTimeoutInSeconds * 1000); confBuilder.setReadTimeout(readTimeoutInSeconds * 1000); confBuilder.setUserAgent(String.format("Pulsar-Java-v%s", getPulsarClientVersion())); - confBuilder.setKeepAliveStrategy((request, httpRequest, httpResponse) -> { - return !httpResponse.getStatus().equals(HttpResponseStatus.INTERNAL_SERVER_ERROR); + confBuilder.setKeepAliveStrategy(new DefaultKeepAliveStrategy() { + @Override + public boolean keepAlive(Request ahcRequest, HttpRequest request, HttpResponse response) { + // Close connection upon a server error or per HTTP spec + return (response.getStatus().code() / 100 != 5) && super.keepAlive(ahcRequest, request, response); + } }); if ("https".equals(url.getProtocol())) {