From 61e05205d7c202606f4ade4cd721c01fe124ac36 Mon Sep 17 00:00:00 2001 From: Shounak kulkarni Date: Thu, 30 Dec 2021 09:39:12 +0530 Subject: [PATCH] async-http-client upgraded Changed updated org to AsyncHttpClient Pumped version to 2.12.3 Refactored pinot-java-client and pinot-jdbc-client as per new version --- pinot-clients/pinot-java-client/pom.xml | 6 ++- .../JsonAsyncHttpPinotClientTransport.java | 48 ++++++++++++++----- pinot-clients/pinot-jdbc-client/pom.xml | 6 ++- .../controller/PinotControllerTransport.java | 21 +++++--- .../response/ControllerResponseFuture.java | 5 +- .../ControllerTenantBrokerResponse.java | 2 +- .../controller/response/SchemaResponse.java | 2 +- .../controller/response/TableResponse.java | 2 +- pom.xml | 4 +- 9 files changed, 68 insertions(+), 28 deletions(-) diff --git a/pinot-clients/pinot-java-client/pom.xml b/pinot-clients/pinot-java-client/pom.xml index 5671e368f5a5..26e0c6f8fb89 100644 --- a/pinot-clients/pinot-java-client/pom.xml +++ b/pinot-clients/pinot-java-client/pom.xml @@ -61,13 +61,17 @@ jackson-databind - com.ning + org.asynchttpclient async-http-client io.netty netty + + io.netty + netty-transport-native-unix-common + diff --git a/pinot-clients/pinot-java-client/src/main/java/org/apache/pinot/client/JsonAsyncHttpPinotClientTransport.java b/pinot-clients/pinot-java-client/src/main/java/org/apache/pinot/client/JsonAsyncHttpPinotClientTransport.java index 98687e9a5781..603d6aad0820 100644 --- a/pinot-clients/pinot-java-client/src/main/java/org/apache/pinot/client/JsonAsyncHttpPinotClientTransport.java +++ b/pinot-clients/pinot-java-client/src/main/java/org/apache/pinot/client/JsonAsyncHttpPinotClientTransport.java @@ -22,9 +22,11 @@ import com.fasterxml.jackson.databind.ObjectReader; import com.fasterxml.jackson.databind.node.JsonNodeFactory; import com.fasterxml.jackson.databind.node.ObjectNode; -import com.ning.http.client.AsyncHttpClient; -import com.ning.http.client.AsyncHttpClientConfig; -import com.ning.http.client.Response; +import io.netty.handler.ssl.ClientAuth; +import io.netty.handler.ssl.JdkSslContext; +import io.netty.handler.ssl.SslContext; +import java.io.IOException; +import java.nio.charset.StandardCharsets; import java.util.HashMap; import java.util.Map; import java.util.concurrent.ExecutionException; @@ -33,6 +35,11 @@ import javax.annotation.Nullable; import javax.net.ssl.SSLContext; import org.apache.pinot.spi.utils.CommonConstants; +import org.asynchttpclient.AsyncHttpClient; +import org.asynchttpclient.BoundRequestBuilder; +import org.asynchttpclient.DefaultAsyncHttpClientConfig.Builder; +import org.asynchttpclient.Dsl; +import org.asynchttpclient.Response; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -52,25 +59,38 @@ public class JsonAsyncHttpPinotClientTransport implements PinotClientTransport { public JsonAsyncHttpPinotClientTransport() { _headers = new HashMap<>(); _scheme = CommonConstants.HTTP_PROTOCOL; - _httpClient = new AsyncHttpClient(); + _httpClient = Dsl.asyncHttpClient(); } public JsonAsyncHttpPinotClientTransport(Map headers, String scheme, - @Nullable SSLContext sslContext) { + @Nullable SSLContext sslContext) { _headers = headers; _scheme = scheme; - AsyncHttpClientConfig.Builder builder = new AsyncHttpClientConfig.Builder(); + Builder builder = Dsl.config(); if (sslContext != null) { - builder.setSSLContext(sslContext); + builder.setSslContext(new JdkSslContext(sslContext, true, ClientAuth.OPTIONAL)); } - _httpClient = new AsyncHttpClient(builder.build()); + _httpClient = Dsl.asyncHttpClient(builder.build()); + } + + public JsonAsyncHttpPinotClientTransport(Map headers, String scheme, + @Nullable SslContext sslContext) { + _headers = headers; + _scheme = scheme; + + Builder builder = Dsl.config(); + if (sslContext != null) { + builder.setSslContext(sslContext); + } + + _httpClient = Dsl.asyncHttpClient(builder.build()); } @Override public BrokerResponse executeQuery(String brokerAddress, String query) - throws PinotClientException { + throws PinotClientException { try { return executeQueryAsync(brokerAddress, query).get(); } catch (Exception e) { @@ -97,7 +117,7 @@ public Future executePinotQueryAsync(String brokerAddress, final url = _scheme + "://" + brokerAddress + "/query"; } - AsyncHttpClient.BoundRequestBuilder requestBuilder = _httpClient.preparePost(url); + BoundRequestBuilder requestBuilder = _httpClient.preparePost(url); if (_headers != null) { _headers.forEach((k, v) -> requestBuilder.addHeader(k, v)); @@ -135,7 +155,11 @@ public void close() if (_httpClient.isClosed()) { throw new PinotClientException("Connection is already closed!"); } - _httpClient.close(); + try { + _httpClient.close(); + } catch (IOException exception) { + throw new PinotClientException("Error while closing connection!"); + } } private static class BrokerResponseFuture implements Future { @@ -185,7 +209,7 @@ public BrokerResponse get(long timeout, TimeUnit unit) "Pinot returned HTTP status " + httpResponse.getStatusCode() + ", expected 200"); } - String responseBody = httpResponse.getResponseBody("UTF-8"); + String responseBody = httpResponse.getResponseBody(StandardCharsets.UTF_8); return BrokerResponse.fromJson(OBJECT_READER.readTree(responseBody)); } catch (Exception e) { throw new ExecutionException(e); diff --git a/pinot-clients/pinot-jdbc-client/pom.xml b/pinot-clients/pinot-jdbc-client/pom.xml index 2bc332a297b7..8a7cbeca1a94 100644 --- a/pinot-clients/pinot-jdbc-client/pom.xml +++ b/pinot-clients/pinot-jdbc-client/pom.xml @@ -75,13 +75,17 @@ jackson-databind - com.ning + org.asynchttpclient async-http-client io.netty netty + + io.netty + netty-transport-native-unix-common + diff --git a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/PinotControllerTransport.java b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/PinotControllerTransport.java index 07dd8c1e9bd0..ac6860ce46a7 100644 --- a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/PinotControllerTransport.java +++ b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/PinotControllerTransport.java @@ -18,8 +18,7 @@ */ package org.apache.pinot.client.controller; -import com.ning.http.client.AsyncHttpClient; -import com.ning.http.client.Response; +import java.io.IOException; import java.util.Map; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; @@ -27,6 +26,10 @@ import org.apache.pinot.client.controller.response.ControllerTenantBrokerResponse; import org.apache.pinot.client.controller.response.SchemaResponse; import org.apache.pinot.client.controller.response.TableResponse; +import org.asynchttpclient.AsyncHttpClient; +import org.asynchttpclient.BoundRequestBuilder; +import org.asynchttpclient.Dsl; +import org.asynchttpclient.Response; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -35,7 +38,7 @@ public class PinotControllerTransport { private static final Logger LOGGER = LoggerFactory.getLogger(PinotControllerTransport.class); - AsyncHttpClient _httpClient = new AsyncHttpClient(); + AsyncHttpClient _httpClient = Dsl.asyncHttpClient(); Map _headers; public PinotControllerTransport() { @@ -48,7 +51,7 @@ public PinotControllerTransport(Map headers) { public TableResponse getAllTables(String controllerAddress) { try { String url = "http://" + controllerAddress + "/tables"; - AsyncHttpClient.BoundRequestBuilder requestBuilder = _httpClient.prepareGet(url); + BoundRequestBuilder requestBuilder = _httpClient.prepareGet(url); if (_headers != null) { _headers.forEach((k, v) -> requestBuilder.addHeader(k, v)); } @@ -66,7 +69,7 @@ public TableResponse getAllTables(String controllerAddress) { public SchemaResponse getTableSchema(String table, String controllerAddress) { try { String url = "http://" + controllerAddress + "/tables/" + table + "/schema"; - AsyncHttpClient.BoundRequestBuilder requestBuilder = _httpClient.prepareGet(url); + BoundRequestBuilder requestBuilder = _httpClient.prepareGet(url); if (_headers != null) { _headers.forEach((k, v) -> requestBuilder.addHeader(k, v)); } @@ -84,7 +87,7 @@ public SchemaResponse getTableSchema(String table, String controllerAddress) { public ControllerTenantBrokerResponse getBrokersFromController(String controllerAddress, String tenant) { try { String url = "http://" + controllerAddress + "/v2/brokers/tenants/" + tenant; - AsyncHttpClient.BoundRequestBuilder requestBuilder = _httpClient.prepareGet(url); + BoundRequestBuilder requestBuilder = _httpClient.prepareGet(url); if (_headers != null) { _headers.forEach((k, v) -> requestBuilder.addHeader(k, v)); } @@ -105,6 +108,10 @@ public void close() if (_httpClient.isClosed()) { throw new PinotClientException("Connection is already closed!"); } - _httpClient.close(); + try { + _httpClient.close(); + } catch (IOException exception) { + throw new PinotClientException("Error while closing connection!"); + } } } diff --git a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerResponseFuture.java b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerResponseFuture.java index 485a4a8eb079..1bc0470d75de 100644 --- a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerResponseFuture.java +++ b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerResponseFuture.java @@ -18,11 +18,12 @@ */ package org.apache.pinot.client.controller.response; -import com.ning.http.client.Response; +import java.nio.charset.StandardCharsets; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import org.apache.pinot.client.PinotClientException; +import org.asynchttpclient.Response; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -74,7 +75,7 @@ public String getStringResponse(long timeout, TimeUnit unit) throw new PinotClientException("Pinot returned HTTP status " + httpResponse.getStatusCode() + ", expected 200"); } - String responseBody = httpResponse.getResponseBody("UTF-8"); + String responseBody = httpResponse.getResponseBody(StandardCharsets.UTF_8); return responseBody; } catch (Exception e) { diff --git a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerTenantBrokerResponse.java b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerTenantBrokerResponse.java index 5872219e5182..4bc3d650089a 100644 --- a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerTenantBrokerResponse.java +++ b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/ControllerTenantBrokerResponse.java @@ -21,13 +21,13 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectReader; -import com.ning.http.client.Response; import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import org.asynchttpclient.Response; public class ControllerTenantBrokerResponse { diff --git a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/SchemaResponse.java b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/SchemaResponse.java index 719bd0825e6d..76e737679e0e 100644 --- a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/SchemaResponse.java +++ b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/SchemaResponse.java @@ -21,11 +21,11 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectReader; -import com.ning.http.client.Response; import java.io.IOException; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import org.asynchttpclient.Response; public class SchemaResponse { diff --git a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/TableResponse.java b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/TableResponse.java index 29d57ab3c178..16dc2399c44f 100644 --- a/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/TableResponse.java +++ b/pinot-clients/pinot-jdbc-client/src/main/java/org/apache/pinot/client/controller/response/TableResponse.java @@ -21,13 +21,13 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectReader; -import com.ning.http.client.Response; import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import org.asynchttpclient.Response; public class TableResponse { diff --git a/pom.xml b/pom.xml index 8adb23e435ff..973f4c245211 100644 --- a/pom.xml +++ b/pom.xml @@ -122,7 +122,7 @@ 0.7 2.10.0 3.5.8 - 1.9.21 + 2.12.3 2.28 2.4.4 1.5.16 @@ -720,7 +720,7 @@ 3.3.4 - com.ning + org.asynchttpclient async-http-client ${async-http-client.version}