如何在Spring Boot应用中通过ElasticSearch异步客户端监控ES服务器健康状态?
Spring Boot 用Elasticsearch异步客户端监控集群健康状态
1. 添加依赖
在Maven的pom.xml中引入Elasticsearch Java异步客户端依赖,注意版本要和你的ES服务器完全一致:
<dependency> <groupId>co.elastic.clients</groupId> <artifactId>elasticsearch-java</artifactId> <version>8.11.3</version> <!-- 替换为你的ES版本 --> </dependency> <!-- 可选:处理HTTPS环境下的SSL证书问题 --> <dependency> <groupId>io.github.hakky54</groupId> <artifactId>sslcontext-kickstart</artifactId> <version>7.1.0</version> </dependency>
如果用Gradle,添加到build.gradle:
implementation 'co.elastic.clients:elasticsearch-java:8.11.3' implementation 'io.github.hakky54:sslcontext-kickstart:7.1.0'
2. 配置异步客户端Bean
创建Spring配置类,初始化Elasticsearch异步客户端,包含节点地址、认证(如果需要)和SSL设置:
import co.elastic.clients.elasticsearch.ElasticsearchAsyncClient; import co.elastic.clients.json.jackson.JacksonJsonpMapper; import co.elastic.clients.transport.ElasticsearchTransport; import co.elastic.clients.transport.rest_client.RestClientTransport; import 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.impl.client.BasicCredentialsProvider; import org.elasticsearch.client.RestClient; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class ElasticsearchConfig { @Bean public ElasticsearchAsyncClient elasticsearchAsyncClient() { // 账号密码认证配置(如果ES开启了安全验证) CredentialsProvider credentialsProvider = new BasicCredentialsProvider(); credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials("你的用户名", "你的密码")); // 构建RestClient RestClient restClient = RestClient.builder( new HttpHost("localhost", 9200, "http")) // 替换为你的ES节点地址 .setHttpClientConfigCallback(httpClientBuilder -> httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider)) .build(); // 创建传输层并初始化异步客户端 ElasticsearchTransport transport = new RestClientTransport(restClient, new JacksonJsonpMapper()); return new ElasticsearchAsyncClient(transport); } }
如果你的ES是HTTPS协议,需要添加SSL证书处理(测试环境可临时忽略证书,生产环境请使用合法证书):
// 在构建RestClient前添加SSL配置 import io.github.hakky54.sslcontext.SSLFactory; import javax.net.ssl.SSLContext; import org.apache.http.conn.ssl.NoopHostnameVerifier; SSLContext sslContext = SSLFactory.builder() .withUnsafeTrustMaterial() .withUnsafeHostnameVerifier() .build() .getSslContext(); // 修改HttpClientConfigCallback部分 httpClientBuilder.setSSLContext(sslContext) .setSSLHostnameVerifier(NoopHostnameVerifier.INSTANCE);
3. 编写健康状态监控服务
创建服务类,调用ES集群健康API获取状态(GREEN/RED/YELLOW):
import co.elastic.clients.elasticsearch.ElasticsearchAsyncClient; import co.elastic.clients.elasticsearch.cluster.HealthResponse; import org.springframework.stereotype.Service; import java.util.concurrent.CompletableFuture; @Service public class EsHealthMonitorService { private final ElasticsearchAsyncClient esAsyncClient; public EsHealthMonitorService(ElasticsearchAsyncClient esAsyncClient) { this.esAsyncClient = esAsyncClient; } // 异步获取集群健康状态 public CompletableFuture<String> getClusterHealthStatus() { return esAsyncClient.cluster().health() .thenApply(HealthResponse::status) .thenApply(Enum::name); } }
4. 测试健康状态查询
可以写一个Controller对外提供接口,或者用单元测试验证:
import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import java.util.concurrent.CompletableFuture; @RestController public class EsHealthController { private final EsHealthMonitorService healthMonitorService; public EsHealthController(EsHealthMonitorService healthMonitorService) { this.healthMonitorService = healthMonitorService; } @GetMapping("/es/health") public CompletableFuture<String> getEsClusterHealth() { return healthMonitorService.getClusterHealthStatus(); } }
关键注意事项
- 版本匹配:Elasticsearch客户端版本必须与服务器版本完全一致,否则会出现兼容性问题。
- 异常处理:实际项目中要捕获
ElasticsearchException等异常,处理连接失败、认证错误等场景。 - 异步特性:异步客户端返回
CompletableFuture,Spring Web会自动处理异步响应,无需手动阻塞线程。
内容的提问来源于stack exchange,提问作者vyshnavi
相关产品推荐
相关产品推荐

