Облачная платформаAdvanced

Подключение к кластеру с использованием Low Level REST Client

Язык статьи: Русский
Показать оригинал
Страница переведена автоматически и может содержать неточности. Рекомендуем сверяться с английской версией.

При запросе и управлении данными в Elasticsearch иногда High Level REST Client может оказаться недостаточным. Например, из‑за встроенных ограничений High Level REST Client может не выполнять сложные пользовательские запросы. В этом случае можно рассмотреть использование Low Level REST Client, который инкапсулирует Elasticsearch API. Достаточно сконструировать необходимые структуры запросов для доступа к кластеру Elasticsearch. Это упрощает работу с кластерами Elasticsearch. Low Level REST Client позволяет настраивать структуру запроса, что более гибко и поддерживает все форматы запросов Elasticsearch, такие как GET, POST, DELETE и HEAD.

Вы можете использовать Low Level REST Client для доступа к кластеру Elasticsearch одним из следующих способов:

  • Создайте Low Level Rest Client напрямую.
  • Создайте High Level Rest Client, а затем вызовите getLowLevelClient() для получения Low Level Rest Client. То есть используйте RestHighLevelClient.getLowLevelClient() для получения Low Level Rest Client.

Как определить, какой метод использовать? Если необходимо выполнять сильно кастомизированные запросы, создайте Low Level Rest Client напрямую. Если вы уже используете High Level Rest Client, можете вызвать метод getLowLevelClient() для получения Low Level Rest Client. Это упрощает ваш код.

Требования

  • Целевой кластер Elasticsearch доступен.
  • Сервер, на котором выполняется Java‑код, может связываться с кластером Elasticsearch.
  • В зависимости от используемого метода настройки сети получите адрес доступа к кластеру. Подробности см. в Obtaining the Cluster Access Address.
  • Java установлена на сервере, версия JDK — 1.8 или новее. Скачайте JDK 1.8 с Java Downloads.
  • Версия Low Level REST Client подтверждена. CSS позволяет подключаться к кластеру Elasticsearch с помощью Java‑клиента более новой версии. Однако для обеспечения лучшей совместимости рекомендуется использовать Java‑клиент той же версии, что и целевой кластер Elasticsearch.

Введение зависимостей

Укажите необходимые Java‑зависимости на сервере, где вы запускаете Java‑код. Объявите версию Apache в режиме Maven.

Замените 7.10.2 на фактическую версию Java‑клиента.

<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>

Подключение к кластеру

Пример кода зависит от настроек режима безопасности целевого кластера Elasticsearch. Выберите соответствующий справочный документ в зависимости от вашего сценария обслуживания.

Table 1 Сценарии доступа к кластеру

Как создаётся Low Level REST Client

Настройки режима безопасности кластера Elasticsearch

Нужно ли загружать сертификат безопасности

Details

Создать Low Level REST Client напрямую

Режим без безопасности

-

Подключение к кластеру в режиме без безопасности с использованием Low Level REST Client

Security mode + HTTP

Security mode + HTTPS

No

Подключение к кластеру в режиме безопасности с использованием Low Level REST Client (без сертификата)

Security mode + HTTPS

Yes

Подключение к кластеру Security-Mode с использованием Low Level REST Client (с сертификатом)

Сначала создайте High Level REST Client, а затем вызовите getLowLevelClient() для получения Low Level REST Client

Режим без безопасности

-

Подключение к кластеру Non-Security Mode с использованием High Level REST Client

Security mode + HTTP

Security mode + HTTPS

No

Подключение к кластеру Security-Mode с использованием High Level REST Client (без сертификата)

Security mode + HTTPS

Yes

Подключение к кластеру Security-Mode с использованием High Level REST Client (с сертификатом)

Подключение к кластеру Non-Security Mode с использованием Low Level REST Client

Используйте Low Level REST Client для подключения к кластеру Elasticsearch, у которого отключён security mode, и выполните запрос, существует ли индекс test. Пример кода приведён ниже:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35import 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);
}
}

Этот фрагмент кода проверяет, существует ли в кластере индекс test. Если возвращён статус 200 (индекс существует) или 404 (индекс не существует), это указывает на то, что соединение с кластером установлено.

Подключение к Security-Mode Cluster с использованием Low Level REST Client (без сертификата)

Используйте Low Level REST Client для подключения к кластеру security-mode Elasticsearch (HTTP или HTTPS) без загрузки сертификата безопасности и выполните запрос, существует ли индекс test. Пример кода приведён ниже:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108import 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 authentication
if (username != null && password != null) {
final CredentialsProvider credentialsProvider = new BasicCredentialsProvider();
credentialsProvider.setCredentials(AuthScope.ANY,
new UsernamePasswordCredentials(username, password));
httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);
}
// set keepalive
httpClientBuilder.setKeepAliveStrategy(((httpResponse, httpContext) -> TimeUnit.MINUTES.toMinutes(10)));
// enable SSL / TLS
SSLContext 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() {
@Override
public void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public 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();
}
}
Table 2 Variables

Parameter

Description

host

IP-адрес для доступа к кластеру. Если указано несколько IP-адресов, разделите их запятой (,).

port

Порт доступа к кластеру. Значение по умолчанию — 9200.

protocol

Протокол соединения, который может быть http или https.

connectTimeout

Тайм‑аут сокет‑соединения (в мс).

connectionRequestTimeout

Тайм‑аут запроса сокет‑соединения (в мс).

socketTimeout

Таймаут запроса сокета (в мс).

username

Имя пользователя для доступа к кластеру.

password

Пароль пользователя.

Этот фрагмент кода проверяет, существует ли индекс test в кластере. Если возвращается 200 (индекс существует) или 404 (индекс не существует), это указывает на то, что кластер подключён.

Подключение к кластеру в режиме Security-Mode с использованием Low Level REST Client (с сертификатом)

Используйте Low Level REST Client для подключения к кластеру Elasticsearch в режиме security-mode, использующему HTTPS с загруженным сертификатом безопасности, и проверьте, существует ли индекс test. Пример кода приведён ниже:

Caution

Чтобы узнать, как получить и загрузить сертификат безопасности, см. Obtaining and Uploading a Security Certificate.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133import 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 authentication
if (username != null && password != null) {
final CredentialsProvider credentialsProvider = new BasicCredentialsProvider();
credentialsProvider.setCredentials(AuthScope.ANY,
new UsernamePasswordCredentials(username, password));
httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);
}
// set keepalive
httpClientBuilder.setKeepAliveStrategy(((httpResponse, httpContext) -> TimeUnit.MINUTES.toMinutes(10)));
// enable SSL / TLS
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 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");
}
@Override
public void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public 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();
}
}
Table 3 Параметры функции

Параметр

Описание

host

IP-адрес для доступа к кластеру. Если указано несколько IP-адресов, разделите их запятой (,).

port

Порт доступа к кластеру. Значение по умолчанию 9200.

protocol

Протокол соединения. Установите этот параметр в https.

connectTimeout

Тайм‑аут сокет‑соединения (в мс).

connectionRequestTimeout

Тайм‑аут запроса сокет‑соединения (в мс).

socketTimeout

Тайм‑аут сокет‑запроса (в мс).

username

Имя пользователя для доступа к кластеру.

password

Пароль пользователя.

certFilePath

Путь для хранения сертификата безопасности.

certPassword

Пароль сертификата безопасности.

Этот фрагмент кода проверяет, существует ли индекс test в кластере. Если возвращён 200 (индекс существует) или 404 (индекс не существует), это указывает на то, что кластер подключён.

Подключение к кластеру без режима безопасности с использованием High Level REST Client

Используйте High Level REST Client для получения Low Level REST Client, вызвав getLowLevelClient(), используйте low-level client для подключения к кластеру Elasticsearch, у которого отключён режим безопасности, и запросите, существует ли индекс test. Пример кода приведён ниже:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37import 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);
}
}

Этот фрагмент кода проверяет, существует ли индекс test в кластере. Если возвращён 200 (индекс существует) или 404 (индекс не существует), это указывает на то, что кластер подключён.

Подключение к кластеру в режиме безопасности с использованием High Level REST Client (без сертификата)

Используйте High Level REST Client для получения Low Level REST Client, вызвав getLowLevelClient(), используйте low-level client для подключения к кластеру Elasticsearch в режиме безопасности, который использует HTTP или HTTPS без загрузки сертификата безопасности, и запросите, существует ли индекс test. Пример кода приведён ниже:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196import 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() {
@Override
public void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public 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;
@Override
public 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 {
@Nullable
private 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}.
*/
@Nullable
CredentialsProvider 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}.
*/
@Override
public HttpAsyncClientBuilder customizeHttpClient(final HttpAsyncClientBuilder httpClientBuilder) {
// enable SSL / TLS
httpClientBuilder.setSSLStrategy(sslStrategy);
// enable user authentication
if (credentialsProvider != null) {
httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);
}
return httpClientBuilder;
}
}
public static class NullHostNameVerifier implements HostnameVerifier {
@Override
public 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();
}
}
Table 4 Переменные

Параметр

Описание

host

IP-адрес для доступа к кластеру. Если указано несколько IP-адресов, разделите их запятой (,).

port

Порт доступа к кластеру. Значение по умолчанию — 9200.

protocol

Протокол соединения, который может быть http или https.

connectTimeout

Тайм‑аут соединения сокета (в мс).

connectionRequestTimeout

Тайм‑аут запроса соединения сокета (в мс).

socketTimeout

Тайм‑аут запроса сокета (в мс).

username

Имя пользователя для доступа к кластеру.

password

Пароль пользователя.

Этот фрагмент кода проверяет, существует ли индекс test в кластере. Если возвращён 200 (индекс существует) или 404 (индекс не существует), это указывает на то, что кластер подключён.

Подключение к Security-Mode Cluster с использованием High Level REST Client (с сертификатом)

Используйте High Level REST Client для получения Low Level REST Client, вызвав getLowLevelClient(), используйте low-level client для подключения к security-mode Elasticsearch cluster, использующему HTTPS с загруженным сертификатом безопасности, и запросите, существует ли индекс test. Пример кода приведён ниже:

Caution

Чтобы узнать, как получить и загрузить сертификат безопасности, см. Obtaining and Uploading a Security Certificate.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162import 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 {
@Nullable
private final CredentialsProvider credentialsProvider;
private final SSLIOSessionStrategy sslStrategy;
SecuredHttpClientConfigCallback(final SSLIOSessionStrategy sslStrategy,
@Nullable final CredentialsProvider credentialsProvider) {
this.sslStrategy = Objects.requireNonNull(sslStrategy);
this.credentialsProvider = credentialsProvider;
}
@Nullable
CredentialsProvider getCredentialsProvider() {
return credentialsProvider;
}
SSLIOSessionStrategy getSSLStrategy() {
return sslStrategy;
}
@Override
public 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");
}
@Override
public void checkClientTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public void checkServerTrusted(X509Certificate[] chain, String authType) throws CertificateException {
}
@Override
public 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();
}
}
Table 5 Параметры функции

Параметр

Описание

host

IP-адрес для доступа к кластеру. Если указано несколько IP-адресов, разделите их запятой (,).

port

Порт доступа к кластеру. Значение по умолчанию — 9200.

protocol

Протокол соединения. Установите этот параметр в https.

connectTimeout

Тайм‑аут сокет‑соединения (в мс).

connectionRequestTimeout

Тайм‑аут запроса сокет‑соединения (в мс).

socketTimeout

Тайм-аут запроса сокета (в мс).

username

Имя пользователя для доступа к кластеру.

password

Пароль пользователя.

certFilePath

Путь для хранения сертификата безопасности.

certPassword

Пароль сертификата безопасности.

Этот фрагмент кода проверяет, существует ли индекс test в кластере. Если возвращается 200 (индекс существует) или 404 (индекс не существует), это указывает на то, что кластер подключён.

Получение и загрузка сертификата безопасности

Для доступа к кластеру Elasticsearch в режиме безопасности, использующему HTTPS, выполните следующие шаги, чтобы при необходимости получить сертификат безопасности и загрузить его в клиент.

  1. Получите сертификат безопасности CloudSearchService.cer.
    1. Войдите в консоль управления CSS.
    2. В панели навигации слева выберите Clusters > Elasticsearch.
    3. В списке кластеров нажмите имя целевого кластера. Отобразится страница информации о кластере.
    4. Нажмите вкладку Overview. В области Network Information нажмите Download Certificate под HTTPS Access.
  2. Преобразуйте сертификат безопасности CloudSearchService.cer. Загрузите загруженный сертификат безопасности на клиент и используйте keytool для преобразования сертификата .cer в сертификат .jks, который может быть прочитан Java.
    • В Linux выполните следующую команду для преобразования сертификата:
      keytool -import -alias newname -keystore ./truststore.jks -file ./CloudSearchService.cer
    • В Windows выполните следующую команду для преобразования сертификата:
      keytool -import -alias newname -keystore .\truststore.jks -file .\CloudSearchService.cer

    В приведённой выше команде newname указывает пользовательское имя сертификата.

    После выполнения этой команды вам будет предложено задать пароль сертификата и подтвердить его. Надёжно сохраните пароль. Он будет использоваться для доступа к кластеру.