You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 09:47:41