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}