When querying and managing data in Elasticsearch, sometimes you may find the High Level REST Client inadequate. For instance, due to inherent limitations, the High Level REST Client may not be able to execute complex custom requests. In this case, you may consider using the Low Level REST Client, which encapsulates Elasticsearch APIs. You only need to construct the required request structures to access an Elasticsearch cluster. This simplifies the process of working with Elasticsearch clusters. The Low Level REST Client allows you to customize the request structure, which is more flexible and supports all the request formats of Elasticsearch, such as GET, POST, DELETE, and HEAD.
You can use the Low Level REST Client to access an Elasticsearch cluster in either of the following ways:
How do I determine when to use which method? If you need to execute highly customized requests, create the Low Level Rest Client directly. If you are already using the High Level Rest Client, you can call the getLowLevelClient() method to obtain the Low Level Rest Client. This simplifies your code.
Introduce the required Java dependencies on the server where you run Java code. Declare the Apache version in Maven mode.
Replace 7.10.2 with the actual Java client version.
<dependency><groupId>org.elasticsearch.client</groupId><artifactId>elasticsearch-rest-high-level-client</artifactId><version>7.10.2</version></dependency><dependency><groupId>org.elasticsearch</groupId><artifactId>elasticsearch</artifactId><version>7.10.2</version></dependency>
The sample code varies depending on the security mode settings of the target Elasticsearch cluster. Select the right reference document based on your service scenario.
How Low Level REST Client Is Created | Elasticsearch Cluster Security-Mode Settings | Whether to Load a Security Certificate | Details |
|---|---|---|---|
Directly create the Low Level REST Client | Non-security mode | - | Connecting to a Non-Security Mode Cluster Using the Low Level REST Client |
Security mode + HTTP Security mode + HTTPS | No | Connecting to a Security-Mode Cluster Using the Low Level REST Client (Without a Certificate) | |
Security mode + HTTPS | Yes | Connecting to a Security-Mode Cluster Using the Low Level REST Client (With a Certificate) | |
Create the High Level REST Client first and then call getLowLevelClient() to obtain the Low Level REST Client | Non-security mode | - | Connecting to a Non-Security Mode Cluster Using the High Level REST Client |
Security mode + HTTP Security mode + HTTPS | No | Connecting to a Security-Mode Cluster Using the High Level REST Client (Without a Certificate) | |
Security mode + HTTPS | Yes | Connecting to a Security-Mode Cluster Using the High Level REST Client (With a Certificate) |
Use the Low Level REST Client to connect to an Elasticsearch cluster for which the security mode is disabled, and query whether the test index exists. The sample code is as follows:
1234567891011121314151617181920212223242526272829303132333435import org.apache.http.HttpHost;import org.elasticsearch.client.Request;import org.elasticsearch.client.Response;import org.elasticsearch.client.RestClient;import org.elasticsearch.client.RestClientBuilder;import java.io.IOException;import java.util.Arrays;import java.util.List;public class Main {public static void main(String[] args) throws IOException {List<String> host = Arrays.asList("{Cluster access address}");RestClientBuilder builder = RestClient.builder(constructHttpHosts(host, 9200, "http"));/***Create the Low Level Rest Client.*/RestClient lowLevelClient = builder.build();/*** Check whether the test index exists. If the index exists, 200 is returned. If the index does not exist, 404 is returned.*/Request request = new Request("HEAD", "/test");Response response = lowLevelClient.performRequest(request);System.out.println(response.getStatusLine().getStatusCode());lowLevelClient.close();}/*** Use the constructHttpHosts function to convert the node IP address list of the host cluster.*/public static HttpHost[] constructHttpHosts(List<String> host, int port, String protocol) {return host.stream().map(p -> new HttpHost(p, port, protocol)).toArray(HttpHost[]::new);}}
This piece of code checks whether the test index exists in the cluster. If 200 (the index exists) or 404 (the index does not exist) is returned, it indicates that the cluster is connected.
Use the Low Level REST Client to connect to a security-mode Elasticsearch cluster (HTTP or HTTPS) without loading a security certificate, and query whether the test index exists. The sample code is as follows:
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108import org.apache.http.HttpHost;import org.apache.http.auth.AuthScope;import org.apache.http.auth.UsernamePasswordCredentials;import org.apache.http.client.CredentialsProvider;import org.apache.http.conn.ssl.NoopHostnameVerifier;import org.apache.http.impl.client.BasicCredentialsProvider;import org.apache.http.nio.conn.ssl.SSLIOSessionStrategy;import org.apache.logging.log4j.LogManager;import org.apache.logging.log4j.Logger;import org.elasticsearch.client.Request;import org.elasticsearch.client.Response;import org.elasticsearch.client.RestClient;import org.elasticsearch.client.RestClientBuilder;import java.io.IOException;import java.security.KeyManagementException;import java.security.NoSuchAlgorithmException;import java.security.SecureRandom;import java.security.cert.CertificateException;import java.security.cert.X509Certificate;import java.util.Arrays;import java.util.List;import java.util.concurrent.TimeUnit;import javax.net.ssl.SSLContext;import javax.net.ssl.TrustManager;import javax.net.ssl.X509TrustManager;public class Main {private static final Logger logger = LogManager.getLogger(Main.class);/*** Create a class for the client. Define the create function.*/public static RestClient create(List<String> host, int port, String protocol, int connectTimeout,int connectionRequestTimeout, int socketTimeout, String username, String password) throws IOException {RestClientBuilder builder = RestClient.builder(constructHttpHosts(host, port, protocol)).setRequestConfigCallback(requestConfig -> requestConfig.setConnectTimeout(connectTimeout).setConnectionRequestTimeout(connectionRequestTimeout).setSocketTimeout(socketTimeout)).setHttpClientConfigCallback(httpClientBuilder -> {// enable user authenticationif (username != null && password != null) {final CredentialsProvider credentialsProvider = new BasicCredentialsProvider();credentialsProvider.setCredentials(AuthScope.ANY,new UsernamePasswordCredentials(username, password));httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);}// set keepalivehttpClientBuilder.setKeepAliveStrategy(((httpResponse, httpContext) -> TimeUnit.MINUTES.toMinutes(10)));// enable SSL / TLSSSLContext sc = null;try {sc = SSLContext.getInstance("SSL");sc.init(null, trustAllCerts, new SecureRandom());} catch (KeyManagementException | NoSuchAlgorithmException e) {e.printStackTrace();}SSLIOSessionStrategy sslStrategy = new SSLIOSessionStrategy(sc, new NoopHostnameVerifier());httpClientBuilder.setSSLStrategy(sslStrategy);return httpClientBuilder;});final RestClient client = builder.build();logger.info("es rest client build success {} ", client);return client;}/*** Use the constructHttpHosts function to convert the node IP address list of the host cluster.*/public static HttpHost[] constructHttpHosts(List<String> host, int port, String protocol) {return host.stream().map(p -> new HttpHost(p, port, protocol)).toArray(HttpHost[]::new);}/*** Configure trustAllCerts to ignore the certificate configuration.*/public static TrustManager[] trustAllCerts = new TrustManager[] {new X509TrustManager() {@Overridepublic void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic X509Certificate[] getAcceptedIssuers() {return null;}}};/*** The following is an example of the main function. Call the create function to create the Low Level REST Client and check whether the test index exists.*/public static void main(String[] args) throws IOException {RestClient lowLevelClient = create(Arrays.asList("{Cluster access address}"), 9200, "http", 1000, 1000, 1000, "username","password");Request request = new Request("HEAD", "/test");Response response = lowLevelClient.performRequest(request);System.out.println(response.getStatusLine().getStatusCode());lowLevelClient.close();}}
Parameter | Description |
|---|---|
host | IP address for accessing the cluster. If there are multiple IP addresses, separate them using a comma (,). |
port | Access port of the cluster. The default value is 9200. |
protocol | Connection protocol, which can be http or https. |
connectTimeout | Socket connection timeout (in ms). |
connectionRequestTimeout | Socket connection request timeout (in ms). |
socketTimeout | Socket request timeout (in ms). |
username | Username for accessing the cluster. |
password | Password of the user. |
This piece of code checks whether the test index exists in the cluster. If 200 (the index exists) or 404 (the index does not exist) is returned, it indicates that the cluster is connected.
Use the Low Level REST Client to connect to a security-mode Elasticsearch cluster that uses HTTPS with a security certificate loaded, and query whether the test index exists. The sample code is as follows:
For how to obtain and upload a security certificate, see Obtaining and Uploading a Security Certificate.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133import org.apache.http.HttpHost;import org.apache.http.auth.AuthScope;import org.apache.http.auth.UsernamePasswordCredentials;import org.apache.http.client.CredentialsProvider;import org.apache.http.conn.ssl.NoopHostnameVerifier;import org.apache.http.impl.client.BasicCredentialsProvider;import org.apache.http.nio.conn.ssl.SSLIOSessionStrategy;import org.apache.logging.log4j.LogManager;import org.apache.logging.log4j.Logger;import org.elasticsearch.client.Request;import org.elasticsearch.client.Response;import org.elasticsearch.client.RestClient;import org.elasticsearch.client.RestClientBuilder;import java.io.File;import java.io.FileInputStream;import java.io.IOException;import java.io.InputStream;import java.security.KeyStore;import java.security.SecureRandom;import java.security.cert.CertificateException;import java.security.cert.X509Certificate;import java.util.Arrays;import java.util.List;import java.util.concurrent.TimeUnit;import javax.net.ssl.SSLContext;import javax.net.ssl.TrustManager;import javax.net.ssl.TrustManagerFactory;import javax.net.ssl.X509TrustManager;public class Main {private static final Logger logger = LogManager.getLogger(Main.class);/*** Create a class for the client. Define the create function.*/public static RestClient create(List<String> host, int port, String protocol, int connectTimeout,int connectionRequestTimeout, int socketTimeout, String username, String password, String certFilePath,String certPassword) throws IOException {RestClientBuilder builder = RestClient.builder(constructHttpHosts(host, port, protocol)).setRequestConfigCallback(requestConfig -> requestConfig.setConnectTimeout(connectTimeout).setConnectionRequestTimeout(connectionRequestTimeout).setSocketTimeout(socketTimeout)).setHttpClientConfigCallback(httpClientBuilder -> {// enable user authenticationif (username != null && password != null) {final CredentialsProvider credentialsProvider = new BasicCredentialsProvider();credentialsProvider.setCredentials(AuthScope.ANY,new UsernamePasswordCredentials(username, password));httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);}// set keepalivehttpClientBuilder.setKeepAliveStrategy(((httpResponse, httpContext) -> TimeUnit.MINUTES.toMinutes(10)));// enable SSL / TLSSSLContext sc = null;try {TrustManager[] tm = {new MyX509TrustManager(certFilePath, certPassword)};sc = SSLContext.getInstance("SSL", "SunJSSE");//You can also use SSLContext sslContext = SSLContext.getInstance("TLSv1.2");sc.init(null, tm, new SecureRandom());} catch (Exception e) {e.printStackTrace();}SSLIOSessionStrategy sslStrategy = new SSLIOSessionStrategy(sc, new NoopHostnameVerifier());httpClientBuilder.setSSLStrategy(sslStrategy);return httpClientBuilder;});final RestClient client = builder.build();logger.info("es rest client build success {} ", client);return client;}/*** Use the constructHttpHosts function to convert the node IP address list of the host cluster.*/public static HttpHost[] constructHttpHosts(List<String> host, int port, String protocol) {return host.stream().map(p -> new HttpHost(p, port, protocol)).toArray(HttpHost[]::new);}public static class MyX509TrustManager implements X509TrustManager {X509TrustManager sunJSSEX509TrustManager;MyX509TrustManager(String certFilePath, String certPassword) throws Exception {File file = new File(certFilePath);if (!file.isFile()) {throw new Exception("Wrong Certification Path");}System.out.println("Loading KeyStore " + file + "...");InputStream in = new FileInputStream(file);KeyStore ks = KeyStore.getInstance("JKS");ks.load(in, certPassword.toCharArray());TrustManagerFactory tmf = TrustManagerFactory.getInstance("SunX509", "SunJSSE");tmf.init(ks);TrustManager[] tms = tmf.getTrustManagers();for (TrustManager tm : tms) {if (tm instanceof X509TrustManager) {sunJSSEX509TrustManager = (X509TrustManager) tm;return;}}throw new Exception("Couldn't initialize");}@Overridepublic void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic X509Certificate[] getAcceptedIssuers() {return new X509Certificate[0];}}/*** The following is an example of the main function. Call the create function to create the Low Level REST Client and check whether the test index exists.*/public static void main(String[] args) throws IOException {RestClient lowLevelClient = create(Arrays.asList("{Cluster access address}"), 9200, "https", 1000, 1000, 1000, "username","password", "certFilePath", "certPassword");Request request = new Request("HEAD", "test");Response response = lowLevelClient.performRequest(request);System.out.println(response.getStatusLine().getStatusCode());lowLevelClient.close();}}
Parameter | Description |
|---|---|
host | IP address for accessing the cluster. If there are multiple IP addresses, separate them using a comma (,). |
port | Access port of the cluster. The default value is 9200. |
protocol | Connection protocol. Set this parameter to https. |
connectTimeout | Socket connection timeout (in ms). |
connectionRequestTimeout | Socket connection request timeout (in ms). |
socketTimeout | Socket request timeout (in ms). |
username | Username for accessing the cluster. |
password | Password of the user. |
certFilePath | Path for storing the security certificate. |
certPassword | Password of the security certificate. |
This piece of code checks whether the test index exists in the cluster. If 200 (the index exists) or 404 (the index does not exist) is returned, it indicates that the cluster is connected.
Use the High Level REST Client to obtain the Low Level REST Client by calling getLowLevelClient(), use the low-level client to connect to an Elasticsearch cluster for which the security mode is disabled, and query whether the test index exists. The sample code is as follows:
12345678910111213141516171819202122232425262728293031323334353637import org.apache.http.HttpHost;import org.elasticsearch.client.Request;import org.elasticsearch.client.Response;import org.elasticsearch.client.RestClient;import org.elasticsearch.client.RestClientBuilder;import org.elasticsearch.client.RestHighLevelClient;import java.io.IOException;import java.util.Arrays;import java.util.List;public class Main {public static void main(String[] args) throws IOException {List<String> host = Arrays.asList("{Cluster access address}");RestClientBuilder builder = RestClient.builder(constructHttpHosts(host, 9200, "http"));final RestHighLevelClient restHighLevelClient = new RestHighLevelClient(builder);/*** Create a High Level Rest Client and then call getLowLevelClient() to obtain the Low Level Rest Client. The code differs from the client creation code only in the following line:*/final RestClient lowLevelClient = restHighLevelClient.getLowLevelClient();/*** Check whether the test index exists. If the index exists, 200 is returned. If the index does not exist, 404 is returned.*/Request request = new Request("HEAD", "/test");Response response = lowLevelClient.performRequest(request);System.out.println(response.getStatusLine().getStatusCode());lowLevelClient.close();}/*** Use the constructHttpHosts function to convert the node IP address list of the host cluster.*/public static HttpHost[] constructHttpHosts(List<String> host, int port, String protocol) {return host.stream().map(p -> new HttpHost(p, port, protocol)).toArray(HttpHost[]::new);}}
This piece of code checks whether the test index exists in the cluster. If 200 (the index exists) or 404 (the index does not exist) is returned, it indicates that the cluster is connected.
Use the High Level REST Client to obtain the Low Level REST Client by calling getLowLevelClient(), use the low-level client to connect to a security-mode Elasticsearch cluster that uses HTTP or HTTPS without loading a security certificate, and query whether the test index exists. The sample code is as follows:
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196import org.apache.http.HttpHost;import org.apache.http.HttpResponse;import org.apache.http.auth.AuthScope;import org.apache.http.auth.UsernamePasswordCredentials;import org.apache.http.client.CredentialsProvider;import org.apache.http.impl.client.BasicCredentialsProvider;import org.apache.http.impl.client.DefaultConnectionKeepAliveStrategy;import org.apache.http.impl.nio.client.HttpAsyncClientBuilder;import org.apache.http.nio.conn.ssl.SSLIOSessionStrategy;import org.apache.http.protocol.HttpContext;import org.apache.logging.log4j.LogManager;import org.apache.logging.log4j.Logger;import org.elasticsearch.client.Request;import org.elasticsearch.client.Response;import org.elasticsearch.client.RestClient;import org.elasticsearch.client.RestClientBuilder;import org.elasticsearch.common.Nullable;import java.io.IOException;import java.security.KeyManagementException;import java.security.NoSuchAlgorithmException;import java.security.SecureRandom;import java.security.cert.CertificateException;import java.security.cert.X509Certificate;import java.util.Arrays;import java.util.List;import java.util.Objects;import java.util.concurrent.TimeUnit;import javax.net.ssl.HostnameVerifier;import javax.net.ssl.SSLContext;import javax.net.ssl.SSLSession;import javax.net.ssl.TrustManager;import javax.net.ssl.X509TrustManager;import org.elasticsearch.client.RestHighLevelClient;public class Main13 {/*** Create a class for the client. Define the create function.*/public static RestHighLevelClient create(List<String> host, int port, String protocol, int connectTimeout, int connectionRequestTimeout, int socketTimeout, String username, String password) throws IOException {final CredentialsProvider credentialsProvider = new BasicCredentialsProvider();credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password));SSLContext sc = null;try {sc = SSLContext.getInstance("SSL");sc.init(null, trustAllCerts, new SecureRandom());} catch (KeyManagementException | NoSuchAlgorithmException e) {e.printStackTrace();}SSLIOSessionStrategy sessionStrategy = new SSLIOSessionStrategy(sc, new NullHostNameVerifier());SecuredHttpClientConfigCallback httpClientConfigCallback = new SecuredHttpClientConfigCallback(sessionStrategy,credentialsProvider);RestClientBuilder builder = RestClient.builder(constructHttpHosts(host, port, protocol)).setRequestConfigCallback(requestConfig -> requestConfig.setConnectTimeout(connectTimeout).setConnectionRequestTimeout(connectionRequestTimeout).setSocketTimeout(socketTimeout)).setHttpClientConfigCallback(httpClientConfigCallback);final RestHighLevelClient client = new RestHighLevelClient(builder);logger.info("es rest client build success {} ", client);return client;}/*** Use the constructHttpHosts function to convert the node IP address list of the host cluster.*/public static HttpHost[] constructHttpHosts(List<String> host, int port, String protocol) {return host.stream().map(p -> new HttpHost(p, port, protocol)).toArray(HttpHost[]::new);}/*** Configure trustAllCerts to ignore the certificate configuration.*/public static TrustManager[] trustAllCerts = new TrustManager[] {new X509TrustManager() {@Overridepublic void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic X509Certificate[] getAcceptedIssuers() {return null;}}};/*** The CustomConnectionKeepAliveStrategy function is used to set the connection keepalive when there are a large number of short connections or when there are not many data requests.*/public static class CustomConnectionKeepAliveStrategy extends DefaultConnectionKeepAliveStrategy {public static final CustomConnectionKeepAliveStrategy INSTANCE = new CustomConnectionKeepAliveStrategy();private CustomConnectionKeepAliveStrategy() {super();}/*** Maximum keepalive time (in minutes)* The default value is 10 minutes. You can set it based on the number of TCP connections in TIME_WAIT state. If there are too many TCP connections, you can increase this value.*/private final long MAX_KEEP_ALIVE_MINUTES = 10;@Overridepublic long getKeepAliveDuration(HttpResponse response, HttpContext context) {long keepAliveDuration = super.getKeepAliveDuration(response, context);// <0 indicates an unlimited keepalive period.// Change the period from unlimited to a default period.if (keepAliveDuration < 0) {return TimeUnit.MINUTES.toMillis(MAX_KEEP_ALIVE_MINUTES);}return keepAliveDuration;}}private static final Logger logger = LogManager.getLogger(Main.class);static class SecuredHttpClientConfigCallback implements RestClientBuilder.HttpClientConfigCallback {@Nullableprivate final CredentialsProvider credentialsProvider;/*** The {@link SSLIOSessionStrategy} for all requests to enable SSL / TLS encryption.*/private final SSLIOSessionStrategy sslStrategy;/*** Create a new {@link SecuredHttpClientConfigCallback}.** @param credentialsProvider The credential provider, if a username/password have been supplied* @param sslStrategy The SSL strategy, if SSL / TLS have been supplied* @throws NullPointerException if {@code sslStrategy} is {@code null}*/SecuredHttpClientConfigCallback(final SSLIOSessionStrategy sslStrategy,@Nullable final CredentialsProvider credentialsProvider) {this.sslStrategy = Objects.requireNonNull(sslStrategy);this.credentialsProvider = credentialsProvider;}/*** Get the {@link CredentialsProvider} that will be added to the HTTP client.** @return Can be {@code null}.*/@NullableCredentialsProvider getCredentialsProvider() {return credentialsProvider;}/*** Get the {@link SSLIOSessionStrategy} that will be added to the HTTP client.** @return Never {@code null}.*/SSLIOSessionStrategy getSSLStrategy() {return sslStrategy;}/*** Sets the {@linkplain HttpAsyncClientBuilder#setDefaultCredentialsProvider(CredentialsProvider) credential provider},** @param httpClientBuilder The client to configure.* @return Always {@code httpClientBuilder}.*/@Overridepublic HttpAsyncClientBuilder customizeHttpClient(final HttpAsyncClientBuilder httpClientBuilder) {// enable SSL / TLShttpClientBuilder.setSSLStrategy(sslStrategy);// enable user authenticationif (credentialsProvider != null) {httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);}return httpClientBuilder;}}public static class NullHostNameVerifier implements HostnameVerifier {@Overridepublic boolean verify(String arg0, SSLSession arg1) {return true;}}/*** The following is an example of the main function. Call the create function to create a high-level client, call the getLowLevelClient() function to obtain a low-level client, and check whether the test index exists.*/public static void main(String[] args) throws IOException {RestHighLevelClient client = create(Arrays.asList("{Cluster access address}") 9200, "http", 1000, 1000, 1000, "username", "password");RestClient lowLevelClient = client.getLowLevelClient();Request request = new Request("HEAD", "test");Response response = lowLevelClient.performRequest(request);System.out.println(response.getStatusLine().getStatusCode());lowLevelClient.close();}}
Parameter | Description |
|---|---|
host | IP address for accessing the cluster. If there are multiple IP addresses, separate them using a comma (,). |
port | Access port of the cluster. The default value is 9200. |
protocol | Connection protocol, which can be http or https. |
connectTimeout | Socket connection timeout (in ms). |
connectionRequestTimeout | Socket connection request timeout (in ms). |
socketTimeout | Socket request timeout (in ms). |
username | Username for accessing the cluster. |
password | Password of the user. |
This piece of code checks whether the test index exists in the cluster. If 200 (the index exists) or 404 (the index does not exist) is returned, it indicates that the cluster is connected.
Use the High Level REST Client to obtain the Low Level REST Client by calling getLowLevelClient(), use the low-level client to connect to a security-mode Elasticsearch cluster that uses HTTPS with a security certificate loaded, and query whether the test index exists. The sample code is as follows:
For how to obtain and upload a security certificate, see Obtaining and Uploading a Security Certificate.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162import org.apache.http.HttpHost;import org.apache.http.auth.AuthScope;import org.apache.http.auth.UsernamePasswordCredentials;import org.apache.http.client.CredentialsProvider;import org.apache.http.conn.ssl.NoopHostnameVerifier;import org.apache.http.impl.client.BasicCredentialsProvider;import org.apache.http.impl.nio.client.HttpAsyncClientBuilder;import org.apache.http.nio.conn.ssl.SSLIOSessionStrategy;import org.apache.logging.log4j.LogManager;import org.apache.logging.log4j.Logger;import org.elasticsearch.action.admin.cluster.health.ClusterHealthRequest;import org.elasticsearch.action.admin.cluster.health.ClusterHealthResponse;import org.elasticsearch.client.Request;import org.elasticsearch.client.RequestOptions;import org.elasticsearch.client.Response;import org.elasticsearch.client.RestClient;import org.elasticsearch.client.RestClientBuilder;import org.elasticsearch.client.RestHighLevelClient;import org.elasticsearch.common.Nullable;import java.io.File;import java.io.FileInputStream;import java.io.IOException;import java.io.InputStream;import java.security.KeyStore;import java.security.SecureRandom;import java.security.cert.CertificateException;import java.security.cert.X509Certificate;import java.util.Arrays;import java.util.List;import java.util.Objects;import javax.net.ssl.SSLContext;import javax.net.ssl.TrustManager;import javax.net.ssl.TrustManagerFactory;import javax.net.ssl.X509TrustManager;public class Main {private static final Logger logger = LogManager.getLogger(Main.class);/*** Create a class for the client. Define the create function.*/public static RestHighLevelClient create(List<String> host, int port, String protocol, int connectTimeout, int connectionRequestTimeout, int socketTimeout, String username, String password, String certFilePath, String certPassword) throws IOException {final CredentialsProvider credentialsProvider = new BasicCredentialsProvider();credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password));SSLContext sc = null;try {TrustManager[] tm = {new MyX509TrustManager(certFilePath, certPassword)};sc = SSLContext.getInstance("SSL", "SunJSSE");//You can also use SSLContext sslContext = SSLContext.getInstance("TLSv1.2");sc.init(null, tm, new SecureRandom());} catch (Exception e) {e.printStackTrace();}SSLIOSessionStrategy sessionStrategy = new SSLIOSessionStrategy(sc, new NoopHostnameVerifier());SecuredHttpClientConfigCallback httpClientConfigCallback = new SecuredHttpClientConfigCallback(sessionStrategy,credentialsProvider);RestClientBuilder builder = RestClient.builder(constructHttpHosts(host, port, protocol)).setRequestConfigCallback(requestConfig -> requestConfig.setConnectTimeout(connectTimeout).setConnectionRequestTimeout(connectionRequestTimeout).setSocketTimeout(socketTimeout)).setHttpClientConfigCallback(httpClientConfigCallback);final RestHighLevelClient client = new RestHighLevelClient(builder);logger.info("es rest client build success {} ", client);ClusterHealthRequest request = new ClusterHealthRequest();ClusterHealthResponse response = client.cluster().health(request, RequestOptions.DEFAULT);logger.info("es rest client health response {} ", response);return client;}/*** Use the constructHttpHosts function to convert the node IP address list of the host cluster.*/public static HttpHost[] constructHttpHosts(List<String> host, int port, String protocol) {return host.stream().map(p -> new HttpHost(p, port, protocol)).toArray(HttpHost[]::new);}static class SecuredHttpClientConfigCallback implements RestClientBuilder.HttpClientConfigCallback {@Nullableprivate final CredentialsProvider credentialsProvider;private final SSLIOSessionStrategy sslStrategy;SecuredHttpClientConfigCallback(final SSLIOSessionStrategy sslStrategy,@Nullable final CredentialsProvider credentialsProvider) {this.sslStrategy = Objects.requireNonNull(sslStrategy);this.credentialsProvider = credentialsProvider;}@NullableCredentialsProvider getCredentialsProvider() {return credentialsProvider;}SSLIOSessionStrategy getSSLStrategy() {return sslStrategy;}@Overridepublic HttpAsyncClientBuilder customizeHttpClient(final HttpAsyncClientBuilder httpClientBuilder) {httpClientBuilder.setSSLStrategy(sslStrategy);if (credentialsProvider != null) {httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);}return httpClientBuilder;}}public static class MyX509TrustManager implements X509TrustManager {X509TrustManager sunJSSEX509TrustManager;MyX509TrustManager(String certFilePath, String certPassword) throws Exception {File file = new File(certFilePath);if (!file.isFile()) {throw new Exception("Wrong Certification Path");}System.out.println("Loading KeyStore " + file + "...");InputStream in = new FileInputStream(file);KeyStore ks = KeyStore.getInstance("JKS");ks.load(in, certPassword.toCharArray());TrustManagerFactory tmf = TrustManagerFactory.getInstance("SunX509", "SunJSSE");tmf.init(ks);TrustManager[] tms = tmf.getTrustManagers();for (TrustManager tm : tms) {if (tm instanceof X509TrustManager) {sunJSSEX509TrustManager = (X509TrustManager) tm;return;}}throw new Exception("Couldn't initialize");}@Overridepublic void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {}@Overridepublic X509Certificate[] getAcceptedIssuers() {return new X509Certificate[0];}}/*** The following is an example of the main function. Call the create function to create a high-level client, call the getLowLevelClient() function to obtain a low-level client, and check whether the test index exists.*/public static void main(String[] args) throws IOException {RestHighLevelClient client = create(Arrays.asList("{Cluster access address}", 9200, "https", 1000, 1000, 1000, "username", "password", "certFilePath", "certPassword");RestClient lowLevelClient = client.getLowLevelClient();Request request = new Request("HEAD", "test");Response response = lowLevelClient.performRequest(request);System.out.println(response.getStatusLine().getStatusCode());lowLevelClient.close();}}
Parameter | Description |
|---|---|
host | IP address for accessing the cluster. If there are multiple IP addresses, separate them using a comma (,). |
port | Access port of the cluster. The default value is 9200. |
protocol | Connection protocol. Set this parameter to https. |
connectTimeout | Socket connection timeout (in ms). |
connectionRequestTimeout | Socket connection request timeout (in ms). |
socketTimeout | Socket request timeout (in ms). |
username | Username for accessing the cluster. |
password | Password of the user. |
certFilePath | Path for storing the security certificate. |
certPassword | Password of the security certificate. |
This piece of code checks whether the test index exists in the cluster. If 200 (the index exists) or 404 (the index does not exist) is returned, it indicates that the cluster is connected.
To access a security-mode Elasticsearch cluster that uses HTTPS, perform the following steps to obtain the security certificate if it is required, and upload it to the client.
keytool -import -alias newname -keystore ./truststore.jks -file ./CloudSearchService.cer
keytool -import -alias newname -keystore .\truststore.jks -file .\CloudSearchService.cer
In the preceding command, newname indicates the user-defined certificate name.
After this command is executed, you will be prompted to set the certificate password and confirm the password. Securely store the password. It will be used for accessing the cluster.