Inhaltsverzeichnis
一、引入依赖jar
二、application.properties配置
三、使用
其它
Heim Java javaLernprogramm So verwenden Sie RestHighLevelClient in springboot2+es7

So verwenden Sie RestHighLevelClient in springboot2+es7

May 15, 2023 pm 10:40 PM
springboot es resthighlevelclient

由于spring和es的集成并不是特别友好,es的高低版本兼容问题、api更新频率高等问题,所以我选择是官网提供的原生Client(RestHighLevelClient),但又不想去关注es的配置类以及和spring的集成配置、jar包冲突等问题,所以使用spring-boot-starter-data-elasticsearch。

一、引入依赖jar

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-elasticsearch</artifactId>
        </dependency>
Nach dem Login kopieren

二、application.properties配置

spring.elasticsearch.rest.uris=http://127.0.0.1:9200,http://127.0.0.1:9201,http://127.0.0.1:9202
spring.elasticsearch.rest.connection-timeout=5s
spring.elasticsearch.rest.read-timeout=30s
logging.level.org.springframework.data.convert.CustomConversions=error
Nach dem Login kopieren

spring-boot-starter-data-elasticsearch中自动装配es的配置类:ElasticsearchRestClientAutoConfiguration、ElasticsearchRestClientProperties。

ElasticsearchRestClientAutoConfiguration:

@ConditionalOnClass({RestHighLevelClient.class})
@ConditionalOnMissingBean({RestClient.class})
@EnableConfigurationProperties({ElasticsearchRestClientProperties.class})
public class ElasticsearchRestClientAutoConfiguration {
 
    @Configuration(
        proxyBeanMethods = false
    )
    @ConditionalOnMissingBean({RestHighLevelClient.class})
    static class RestHighLevelClientConfiguration {
        RestHighLevelClientConfiguration() {
        }
 
        @Bean
        RestHighLevelClient elasticsearchRestHighLevelClient(RestClientBuilder restClientBuilder) {
            return new RestHighLevelClient(restClientBuilder);
        }
    }
 
    @Configuration(
        proxyBeanMethods = false
    )
    @ConditionalOnMissingBean({RestClientBuilder.class})
    static class RestClientBuilderConfiguration {
        RestClientBuilderConfiguration() {
        }
 
        @Bean
        RestClientBuilderCustomizer defaultRestClientBuilderCustomizer(ElasticsearchRestClientProperties properties) {
            return new ElasticsearchRestClientAutoConfiguration.DefaultRestClientBuilderCustomizer(properties);
        }
 
        @Bean
        RestClientBuilder elasticsearchRestClientBuilder(ElasticsearchRestClientProperties properties, ObjectProvider<RestClientBuilderCustomizer> builderCustomizers) {
            HttpHost[] hosts = (HttpHost[])properties.getUris().stream().map(this::createHttpHost).toArray((x$0) -> {
                return new HttpHost[x$0];
            });
            RestClientBuilder builder = RestClient.builder(hosts);
            builder.setHttpClientConfigCallback((httpClientBuilder) -> {
                builderCustomizers.orderedStream().forEach((customizer) -> {
                    customizer.customize(httpClientBuilder);
                });
                return httpClientBuilder;
            });
            builder.setRequestConfigCallback((requestConfigBuilder) -> {
                builderCustomizers.orderedStream().forEach((customizer) -> {
                    customizer.customize(requestConfigBuilder);
                });
                return requestConfigBuilder;
            });
            builderCustomizers.orderedStream().forEach((customizer) -> {
                customizer.customize(builder);
            });
            return builder;
        }
 
        private HttpHost createHttpHost(String uri) {
            try {
                return this.createHttpHost(URI.create(uri));
            } catch (IllegalArgumentException var3) {
                return HttpHost.create(uri);
            }
        }
 
        private HttpHost createHttpHost(URI uri) {
            if (!StringUtils.hasLength(uri.getUserInfo())) {
                return HttpHost.create(uri.toString());
            } else {
                try {
                    return HttpHost.create((new URI(uri.getScheme(), (String)null, uri.getHost(), uri.getPort(), uri.getPath(), uri.getQuery(), uri.getFragment())).toString());
                } catch (URISyntaxException var3) {
                    throw new IllegalStateException(var3);
                }
            }
        }
    }
}
Nach dem Login kopieren

ElasticsearchRestClientProperties:

@ConfigurationProperties(
    prefix = "spring.elasticsearch.rest"
)
public class ElasticsearchRestClientProperties {
    private List<String> uris = new ArrayList(Collections.singletonList("http://localhost:9200"));
    private String username;
    private String password;
    private Duration connectionTimeout = Duration.ofSeconds(1L);
    private Duration readTimeout = Duration.ofSeconds(30L);
 
    public ElasticsearchRestClientProperties() {
    }
 
    public List<String> getUris() {
        return this.uris;
    }
 
    public void setUris(List<String> uris) {
        this.uris = uris;
    }
 
    public String getUsername() {
        return this.username;
    }
 
    public void setUsername(String username) {
        this.username = username;
    }
 
    public String getPassword() {
        return this.password;
    }
 
    public void setPassword(String password) {
        this.password = password;
    }
 
    public Duration getConnectionTimeout() {
        return this.connectionTimeout;
    }
 
    public void setConnectionTimeout(Duration connectionTimeout) {
        this.connectionTimeout = connectionTimeout;
    }
 
    public Duration getReadTimeout() {
        return this.readTimeout;
    }
 
    public void setReadTimeout(Duration readTimeout) {
        this.readTimeout = readTimeout;
    }
}
Nach dem Login kopieren

三、使用

ES基本操作持久层

/**
 * es持久层
 *
 * @author yangzihe
 * @date 2022/1/24
 */
@Repository
@Slf4j
public class EsRepository {
 
    @Resource
    private RestHighLevelClient highLevelClient;
 
    /**
     * 判断索引是否存在
     */
    public boolean existIndex(String index) {
        try {
            return highLevelClient.indices().exists(new GetIndexRequest(index), RequestOptions.DEFAULT);
        } catch (IOException e) {
            log.error("es持久层异常!index={}", index, e);
        }
        return Boolean.FALSE;
    }
 
    /**
     * 创建索引
     */
    public boolean createIndex(String index, Map<String, Object> columnMap) {
        if (existIndex(index)) {
            return Boolean.FALSE;
        }
 
        CreateIndexRequest request = new CreateIndexRequest(index);
        if (columnMap != null && columnMap.size() > 0) {
            Map<String, Object> source = new HashMap<>();
            source.put("properties", columnMap);
            request.mapping(source);
        }
        try {
            highLevelClient.indices().create(request, RequestOptions.DEFAULT);
            return Boolean.TRUE;
        } catch (IOException e) {
            log.error("es持久层异常!index={}, columnMap={}", index, columnMap, e);
        }
        return Boolean.FALSE;
    }
 
    /**
     * 删除索引
     */
    public boolean deleteIndex(String index) {
        try {
            if (existIndex(index)) {
                AcknowledgedResponse response = highLevelClient.indices().delete(new DeleteIndexRequest(index), RequestOptions.DEFAULT);
                return response.isAcknowledged();
            }
        } catch (Exception e) {
            log.error("es持久层异常!index={}", index, e);
        }
        return Boolean.FALSE;
    }
 
    /**
     * 数据新增
     */
    public boolean insert(String index, String jsonString) {
        IndexRequest indexRequest = new IndexRequest(index);
 
        indexRequest.id(new Snowflake().nextIdStr());
        indexRequest.source(jsonString, XContentType.JSON);
 
        try {
            log.info("indexRequest={}", indexRequest);
            IndexResponse indexResponse = highLevelClient.index(indexRequest, RequestOptions.DEFAULT);
            log.info("indexResponse={}", indexResponse);
            return Boolean.TRUE;
        } catch (IOException e) {
            log.error("es持久层异常!index={}, jsonString={}", index, jsonString, e);
        }
        return Boolean.FALSE;
    }
 
    /**
     * 数据更新,可以直接修改索引结构
     */
    public boolean update(String index, Map<String, Object> dataMap) {
        UpdateRequest updateRequest = new UpdateRequest(index, dataMap.remove("id").toString());
        updateRequest.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
        updateRequest.doc(dataMap);
        try {
            highLevelClient.update(updateRequest, RequestOptions.DEFAULT);
        } catch (IOException e) {
            log.error("es持久层异常!index={}, dataMap={}", index, dataMap, e);
            return Boolean.FALSE;
        }
        return Boolean.TRUE;
    }
 
    /**
     * 删除数据
     */
    public boolean delete(String index, String id) {
        DeleteRequest deleteRequest = new DeleteRequest(index, id);
        try {
            highLevelClient.delete(deleteRequest, RequestOptions.DEFAULT);
        } catch (IOException e) {
            log.error("es持久层异常!index={}, id={}", index, id, e);
            return Boolean.FALSE;
        }
        return Boolean.TRUE;
    }
 
}
Nach dem Login kopieren

ES查询持久层

/**
 * es查询持久层
 *
 * @author yangzihe
 * @date 2022/1/25
 */
@Repository
@Slf4j
public class EsSearchRepository {
 
    @Resource
    private RestHighLevelClient highLevelClient;
 
    /**
     * 分页查询
     *
     * @param queryPO 分页查询对象
     *
     * @return 分页查询结果
     */
    public EsQueryRespPO<Map<String, Object>> searchPage(EsQueryReqPO queryPO) {
        // 默认分页参数设置
        if (queryPO.getPageNum() == null) {
            queryPO.setPageNum(1);
        }
        if (queryPO.getPageSize() == null) {
            queryPO.setPageSize(10);
        }
 
        // 设置索引
        SearchRequest searchRequest = new SearchRequest(queryPO.getIndex());
 
        // 封装查询源对象
        SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
        searchRequest.source(sourceBuilder);
 
        // 查询条件
        sourceBuilder.query(queryPO.getQuery());
 
        // 排序字段
        if (StringUtils.isNotBlank(queryPO.getSortField()) && queryPO.getSort() != null) {
            FieldSortBuilder order = new FieldSortBuilder(queryPO.getSortField()).order(queryPO.getSort());
            sourceBuilder.sort(order);
        }
 
        // 开始行数,默认0
        sourceBuilder.from((queryPO.getPageNum() - 1) * queryPO.getPageSize());
        // 页大小,默认10
        sourceBuilder.size(queryPO.getPageSize());
 
        // 查询结果
        SearchResponse searchResponse = null;
        try {
            // log.info("es查询请求对象:searchRequest={}", searchRequest);
            log.info("es查询请求对象source:sourceBuilder={}", searchRequest.source());
            // 执行搜索
            searchResponse = highLevelClient.search(searchRequest, RequestOptions.DEFAULT);
            log.info("es查询响应结果:searchResponse={}", searchResponse);
        } catch (IOException e) {
            log.error("es查询,IO异常!searchRequest={}", searchRequest, e);
            // 异常处理
            return EsQueryRespPO.error("es查询,IO异常!");
        }
 
        if (RestStatus.OK.equals(searchResponse.status())) {
            // 解析对象
            SearchHit[] hits = searchResponse.getHits().getHits();
            // 获取source
            List<Map<String, Object>> sourceList = Arrays.stream(hits).map(SearchHit::getSourceAsMap).collect(Collectors.toList());
            long totalHits = searchResponse.getHits().getTotalHits().value;
            return EsQueryRespPO.success(sourceList, queryPO.getPageNum(), queryPO.getPageSize(), totalHits);
        } else {
            log.error("es查询返回的状态码异常!searchResponse.status={}, searchRequest={}", searchResponse.status(), searchRequest);
            return EsQueryRespPO.error("es查询返回的状态码异常!");
        }
    }
 
    /**
     * 聚合的分页查询
     *
     * @param queryPO 查询请求对象
     *
     * @return 聚合分页查询结果
     */
    public EsQueryRespPO<AggregationBucketPO> searchAggregation(EsQueryReqPO queryPO) {
        // 设置索引
        SearchRequest searchRequest = new SearchRequest(queryPO.getIndex());
 
        // 封装查询源对象
        SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
        searchRequest.source(sourceBuilder);
 
        // 查询条件
        sourceBuilder.query(queryPO.getQuery());
 
        // 排序字段
        if (StringUtils.isNotBlank(queryPO.getSortField()) && queryPO.getSort() != null) {
            FieldSortBuilder order = new FieldSortBuilder(queryPO.getSortField()).order(queryPO.getSort());
            sourceBuilder.sort(order);
        }
 
        // 页大小0,只返回聚合结果
        sourceBuilder.size(0);
 
        // 设置聚合查询,可以设置多个聚合查询条件,只要聚合查询命名不同就行
        // 聚合分组条件, group by
        sourceBuilder.aggregation(queryPO.getTermsAggregation());
        // 聚合统计条件, count分组后的数据,计算分组后的总大小
        sourceBuilder.aggregation(queryPO.getCardinalityAggregation());
 
        // 查询结果
        SearchResponse searchResponse = null;
        try {
            // log.info("es查询请求对象:searchRequest={}", searchRequest);
            log.info("es查询请求对象source:sourceBuilder={}", searchRequest.source());
            // 执行搜索
            searchResponse = highLevelClient.search(searchRequest, RequestOptions.DEFAULT);
            log.info("es查询响应结果:searchResponse={}", searchResponse);
        } catch (IOException e) {
            log.error("es查询,IO异常!searchRequest={}", searchRequest, e);
            return EsQueryRespPO.error("es查询,IO异常!");
        }
 
        if (RestStatus.OK.equals(searchResponse.status())) {
            // 解析对象
            Aggregations aggregations = searchResponse.getAggregations();
 
            long docTotal = searchResponse.getHits().getTotalHits().value;
 
            // 遍历terms聚合结果
            Terms terms = aggregations.get(queryPO.getTermsAggregation().getName());
            List<AggregationBucketPO> bucketList = terms.getBuckets().stream().map(bucket -> {
                String key = bucket.getKeyAsString();
                long docCount = bucket.getDocCount();
                return new AggregationBucketPO(key, docCount, docTotal);
            }).collect(Collectors.toList());
 
            // 总数量
            Cardinality cardinality = aggregations.get(queryPO.getCardinalityAggregation().getName());
            long totalHits = cardinality.getValue();
 
            return EsQueryRespPO.success(bucketList, queryPO.getPageNum(), queryPO.getPageSize(), totalHits);
        } else {
            log.error("es查询返回的状态码异常!searchResponse.status={}, searchRequest={}", searchResponse.status(), searchRequest);
            return EsQueryRespPO.error("es查询返回的状态码异常!");
        }
    }
}
Nach dem Login kopieren

其中,EsQueryReqPO、EsQueryRespPO、AggregationBucketPO等类如下:

/**
 * es查询请求对象
 */
@Data
public class EsQueryReqPO {
 
    /**
     * 索引名
     */
    String[] index;
 
    /**
     * 查询条件
     */
    QueryBuilder query;
 
    /**
     * 排序字段
     */
    String sortField;
 
    /**
     * 排序方式 SortOrder.ASC、SortOrder.DESC
     */
    SortOrder sort;
 
    /**
     * 页数
     */
    private Integer pageNum;
 
    /**
     * 页大小
     */
    private Integer pageSize;
 
    /**
     * 聚合分组条件, group by
     */
    private TermsAggregationBuilder termsAggregation;
 
    /**
     * 聚合统计条件, count分组后的数据
     */
    private CardinalityAggregationBuilder cardinalityAggregation;
}
/**
 * es分页响应对象
 *
 * @author yangzihe
 * @date 2022/1/25
 */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class EsQueryRespPO<T> {
 
    /**
     * 是否成功
     */
    private Boolean success;
 
    /**
     * 信息
     */
    private String message;
 
    /**
     * 页数
     */
    private Integer pageNum;
 
    /**
     * 页大小
     */
    private Integer pageSize;
 
    /**
     * 总大小
     */
    private Long totalSize;
 
    /**
     * 数据
     */
    private List<T> sourceList;
 
    public static <T> EsQueryRespPO<T> success(List<T> sourceList, Integer pageNum, Integer pageSize,
                                               Long totalSize) {
        EsQueryRespPO<T> esQueryRespPO = new EsQueryRespPO<>();
        esQueryRespPO.setSuccess(true);
        esQueryRespPO.setSourceList(sourceList);
        esQueryRespPO.setPageNum(pageNum);
        esQueryRespPO.setPageSize(pageSize);
        esQueryRespPO.setTotalSize(totalSize);
        return esQueryRespPO;
    }
 
    public static EsQueryRespPO error() {
        EsQueryRespPO esQueryRespPO = new EsQueryRespPO();
        esQueryRespPO.setSuccess(false);
        esQueryRespPO.setMessage("es查询异常");
        return esQueryRespPO;
    }
 
    public static EsQueryRespPO error(String message) {
        EsQueryRespPO esQueryRespPO = new EsQueryRespPO();
        esQueryRespPO.setSuccess(false);
        esQueryRespPO.setMessage(message);
        return esQueryRespPO;
    }
 
}
Nach dem Login kopieren
/**
 * 聚合桶对象
 *
 * @author yangzihe
 * @date 2022/1/26
 */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class AggregationBucketPO {
 
    /**
     * 聚合Bucket的key名
     */
    private String key;
 
    /**
     * 聚合Bucket的文档数量
     */
    private Long docCount;
 
    /**
     * 文档总数量
     */
    private Long docTotal;
 
}
Nach dem Login kopieren

ES多级(二级)聚合分桶查询

import com.yy.armor.manager.common.exception.EsException;
import com.yy.armor.manager.persist.es.po.AggregationBucketPO;
import com.yy.armor.manager.persist.es.po.EsMultiAggregationReqPO;
 
import java.io.IOException;
import java.util.List;
import java.util.stream.Collectors;
import javax.annotation.Resource;
 
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
import org.apache.pulsar.shade.org.apache.commons.compress.utils.Lists;
import org.elasticsearch.action.search.SearchRequest;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.rest.RestStatus;
import org.elasticsearch.search.aggregations.AggregationBuilders;
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
import org.elasticsearch.search.aggregations.bucket.terms.TermsAggregationBuilder;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.sort.FieldSortBuilder;
import org.springframework.stereotype.Repository;
 
 
@Repository
@Slf4j
public class EsSearchRepository {
 
    @Resource
    private RestHighLevelClient highLevelClient;
 
 
    /**
     * 多级聚合查询(二级聚合)
     *
     * @param reqPO 查询请求对象
     *
     * @return 聚合查询结果
     */
    public List<AggregationBucketPO> searchMultiAggregation(EsMultiAggregationReqPO reqPO) {
        // 设置索引
        SearchRequest searchRequest = new SearchRequest(reqPO.getIndex());
 
        // 封装查询源对象
        SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
        searchRequest.source(sourceBuilder);
 
        // 查询条件
        sourceBuilder.query(reqPO.getQuery());
 
        // 排序字段
        if (StringUtils.isNotBlank(reqPO.getSortField()) && reqPO.getSort() != null) {
            FieldSortBuilder order = new FieldSortBuilder(reqPO.getSortField()).order(reqPO.getSort());
            sourceBuilder.sort(order);
        }
 
        // 页大小0,只返回聚合结果
        sourceBuilder.size(0);
 
        // 聚合分桶。创建terms桶聚合,聚合名字=terms_by_XXX, 字段=XXX
        TermsAggregationBuilder termsAggregationBuilder = AggregationBuilders.terms("terms_by_" + reqPO.getField()).field(reqPO.getField());
        if (reqPO.getFieldSize() != null) {
            termsAggregationBuilder.size(reqPO.getFieldSize());
        }
        // 二级聚合分桶
        TermsAggregationBuilder subTermsAggregationBuilder = AggregationBuilders.terms("terms_by_" + reqPO.getSubField()).field(reqPO.getSubField());
        if (reqPO.getSubFieldSize() != null) {
            subTermsAggregationBuilder.size(reqPO.getSubFieldSize());
        }
        termsAggregationBuilder.subAggregation(subTermsAggregationBuilder);
 
        // 聚合分组条件
        sourceBuilder.aggregation(termsAggregationBuilder);
 
        // 查询结果
        SearchResponse searchResponse = null;
        try {
            log.info("es查询请求对象source:sourceBuilder={}", searchRequest.source());
            // 执行搜索
            searchResponse = highLevelClient.search(searchRequest, RequestOptions.DEFAULT);
            log.info("es查询响应结果:searchResponse={}", searchResponse);
        } catch (IOException e) {
            log.error("es查询,IO异常!searchRequest={}", searchRequest, e);
            throw new EsException("es查询,IO异常!");
        }
 
        if (RestStatus.OK.equals(searchResponse.status())) {
            // 遍历terms聚合结果
            Terms terms = searchResponse.getAggregations().get(termsAggregationBuilder.getName());
            List<AggregationBucketPO> bucketList = terms.getBuckets().stream().map(bucket -> {
                // 一级聚合分桶的数据
                String key = bucket.getKeyAsString();
                long docCount = bucket.getDocCount();
                // 二级聚合分桶的数据
                Terms subTerms = bucket.getAggregations().get(subTermsAggregationBuilder.getName());
                List<AggregationBucketPO> subBucketList = convertTerms(subTerms);
                return new AggregationBucketPO(key, docCount, subBucketList);
            }).collect(Collectors.toList());
 
            return bucketList;
        } else {
            log.error("es查询返回的状态码异常!searchResponse.status={}, searchRequest={}", searchResponse.status(), searchRequest);
            throw new EsException("es查询返回的状态码异常!");
        }
    }
 
    private List<AggregationBucketPO> convertTerms(Terms terms) {
        if (CollectionUtils.isEmpty(terms.getBuckets())) {
            return Lists.newArrayList();
        }
 
        return terms.getBuckets().stream().map(bucket -> {
            String key = bucket.getKeyAsString();
            long docCount = bucket.getDocCount();
            return new AggregationBucketPO(key, docCount);
        }).collect(Collectors.toList());
    }
}
Nach dem Login kopieren

其中,EsMultiAggregationReqPO、AggregationBucketPO类如下:

@Data
public class EsMultiAggregationReqPO {
    
    /**
     * 索引名
     */
    String[] index;
 
    /**
     * 查询条件
     */
    QueryBuilder query;
 
    /**
     * 聚合分桶字段
     */
    private String field;
 
    /**
     * 二级聚合分桶字段
     */
    private String subField;
 
    /**
     * 聚合分桶大小,非必传
     */
    private Integer fieldSize;
 
    /**
     * 二级聚合分桶大小,非必传
     */
    private Integer subFieldSize;
 
    /**
     * 排序字段,非必传
     */
    String sortField;
 
    /**
     * 排序方式 SortOrder.ASC、SortOrder.DESC,非必传
     */
    SortOrder sort;
}
Nach dem Login kopieren
@Data
@NoArgsConstructor
@AllArgsConstructor
public class AggregationBucketPO {
    /**
     * 聚合Bucket的key名
     */
    private String key;
 
    /**
     * 聚合Bucket的文档数量
     */
    private Long docCount;
 
    /**
     * 子桶集合
     */
    private List<AggregationBucketPO> subBucketList;
 
    public AggregationBucketPO(String key, Long docCount) {
        this.key = key;
        this.docCount = docCount;
    }
}
Nach dem Login kopieren

二级聚合分桶测试代码

    @PostConstruct
    private void init() {
        // 查询对象的封装
        EsMultiAggregationReqPO reqPO = new EsMultiAggregationReqPO();
 
        reqPO.setIndex(new String[]{"test_log"});
 
        List<Long> ids = Lists.newArrayList();
        ids.add(140L);
        ids.add(141L);
        ids.add(142L);
        QueryBuilder queryBuilder4 = QueryBuilders.termsQuery("eventId", ids);
        BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery().must(queryBuilder4);
        reqPO.setQuery(queryBuilder);
 
        reqPO.setField("eventId");
        reqPO.setFieldSize(9999);
        reqPO.setSubField("riskFlag");
 
        // 执行查询
        List<AggregationBucketPO> esQueryRespPO = searchMultiAggregation(reqPO);
        System.out.println("esQueryRespPO=" + esQueryRespPO);
    }
Nach dem Login kopieren

其它

如果没有用spring-boot-starter-data-elasticsearch来自动注入es组件,那么需要自己做es client的注入,es配置类如下:

/**
 * @author yangzihe
 * @date 2022/1/25
 */
@Configuration
public class EsClientConfig {
 
    @Value("${spring.elasticsearch.rest.uris}")
    private List<String> uris;
 
    @Bean
    public RestHighLevelClient restHighLevelClient() {
        List<HttpHost> httpHostList = uris.stream().map(HttpHost::create).collect(Collectors.toList());
        HttpHost[] httpHost = new HttpHost[uris.size()];
        httpHostList.toArray(httpHost);
        RestClientBuilder clientBuilder = RestClient.builder(httpHost);
        return new RestHighLevelClient(clientBuilder);
    }
 
}
Nach dem Login kopieren

Snowflake是hutool包里的,导包:

        <!-- https://mvnrepository.com/artifact/cn.hutool/hutool-all -->
        <dependency>
            <groupId>cn.hutool</groupId>
            <artifactId>hutool-all</artifactId>
            <version>5.7.14</version>
        </dependency>
Nach dem Login kopieren

聚合查询的测试用例:

@RunWith(SpringJUnit4ClassRunner.class)
@SpringBootTest(classes = StartApplication.class)
public class EsTest {
 
    @Resource
    private EsSearchRepository esSearchRepository;
 
    @Test
    public void testSearchAggregation() {
        // 查询对象的封装
        EsQueryReqPO queryPO = new EsQueryReqPO();
        queryPO.setIndex(new String[]{"yzh2", "yzh3"});
        queryPO.setPageNum(1);
        queryPO.setPageSize(10);
 
        // 时间戳范围
        QueryBuilder queryBuilder1 = QueryBuilders.rangeQuery("timestamp")
            .from(System.currentTimeMillis() - 1000)
            .to(System.currentTimeMillis());
        // 登录标识
        QueryBuilder queryBuilder2 = QueryBuilders.termQuery("name", "yang");
        BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery().must(queryBuilder1).must(queryBuilder2);
        queryPO.setQuery(queryBuilder);
 
        // 根据userName分组。创建terms桶聚合,聚合名字=terms_by_userName, 字段=payload.userName.keyword
        TermsAggregationBuilder termsAggregationBuilder = AggregationBuilders
            .terms("terms_by_userName").field("payload.userName.keyword");
        termsAggregationBuilder.size(queryPO.getPageSize() * queryPO.getPageNum());
        termsAggregationBuilder.subAggregation(new BucketSortPipelineAggregationBuilder("bucket_field", null)
            .from((queryPO.getPageNum() - 1) * queryPO.getPageSize()).size(queryPO.getPageSize()));
        queryPO.setTermsAggregation(termsAggregationBuilder);
 
        // 根据userName聚合count统计. cardinality名=count_userName, 字段=payload.userName.keyword
        CardinalityAggregationBuilder cardinalityAggregationBuilder = AggregationBuilders
            .cardinality("count_userName").field("payload.userName.keyword");
        queryPO.setCardinalityAggregation(cardinalityAggregationBuilder);
 
        // 执行查询
        EsQueryRespPO<AggregationBucketPO> esQueryRespPO = esSearchRepository.searchAggregation(queryPO);
    }
}
Nach dem Login kopieren

Das obige ist der detaillierte Inhalt vonSo verwenden Sie RestHighLevelClient in springboot2+es7. Für weitere Informationen folgen Sie bitte anderen verwandten Artikeln auf der PHP chinesischen Website!

Erklärung dieser Website
Der Inhalt dieses Artikels wird freiwillig von Internetnutzern beigesteuert und das Urheberrecht liegt beim ursprünglichen Autor. Diese Website übernimmt keine entsprechende rechtliche Verantwortung. Wenn Sie Inhalte finden, bei denen der Verdacht eines Plagiats oder einer Rechtsverletzung besteht, wenden Sie sich bitte an admin@php.cn

Heiße KI -Werkzeuge

Undresser.AI Undress

Undresser.AI Undress

KI-gestützte App zum Erstellen realistischer Aktfotos

AI Clothes Remover

AI Clothes Remover

Online-KI-Tool zum Entfernen von Kleidung aus Fotos.

Undress AI Tool

Undress AI Tool

Ausziehbilder kostenlos

Clothoff.io

Clothoff.io

KI-Kleiderentferner

AI Hentai Generator

AI Hentai Generator

Erstellen Sie kostenlos Ai Hentai.

Heißer Artikel

R.E.P.O. Energiekristalle erklärten und was sie tun (gelber Kristall)
1 Monate vor By 尊渡假赌尊渡假赌尊渡假赌
R.E.P.O. Beste grafische Einstellungen
1 Monate vor By 尊渡假赌尊渡假赌尊渡假赌
R.E.P.O. So reparieren Sie Audio, wenn Sie niemanden hören können
1 Monate vor By 尊渡假赌尊渡假赌尊渡假赌
R.E.P.O. Chat -Befehle und wie man sie benutzt
1 Monate vor By 尊渡假赌尊渡假赌尊渡假赌

Heiße Werkzeuge

Notepad++7.3.1

Notepad++7.3.1

Einfach zu bedienender und kostenloser Code-Editor

SublimeText3 chinesische Version

SublimeText3 chinesische Version

Chinesische Version, sehr einfach zu bedienen

Senden Sie Studio 13.0.1

Senden Sie Studio 13.0.1

Leistungsstarke integrierte PHP-Entwicklungsumgebung

Dreamweaver CS6

Dreamweaver CS6

Visuelle Webentwicklungstools

SublimeText3 Mac-Version

SublimeText3 Mac-Version

Codebearbeitungssoftware auf Gottesniveau (SublimeText3)

Wie Springboot Jasypt integriert, um die Verschlüsselung von Konfigurationsdateien zu implementieren Wie Springboot Jasypt integriert, um die Verschlüsselung von Konfigurationsdateien zu implementieren Jun 01, 2023 am 08:55 AM

Einführung in Jasypt Jasypt ist eine Java-Bibliothek, die es einem Entwickler ermöglicht, seinem Projekt mit minimalem Aufwand grundlegende Verschlüsselungsfunktionen hinzuzufügen und kein tiefes Verständnis der Funktionsweise der Verschlüsselung erfordert. standardbasierte Verschlüsselungstechnologie. Passwörter, Text, Zahlen, Binärdateien verschlüsseln ... Geeignet für die Integration in Spring-basierte Anwendungen, offene API, zur Verwendung mit jedem JCE-Anbieter ... Fügen Sie die folgende Abhängigkeit hinzu: com.github.ulisesbocchiojasypt-spring-boot-starter2 Die Vorteile von Jasypt schützen unsere Systemsicherheit. Selbst wenn der Code durchgesickert ist, kann die Datenquelle garantiert werden.

Wie SpringBoot Redisson integriert, um eine Verzögerungswarteschlange zu implementieren Wie SpringBoot Redisson integriert, um eine Verzögerungswarteschlange zu implementieren May 30, 2023 pm 02:40 PM

Nutzungsszenario 1. Die Bestellung wurde erfolgreich aufgegeben, die Zahlung erfolgte jedoch nicht innerhalb von 30 Minuten. Die Zahlung ist abgelaufen und die Bestellung wurde automatisch storniert. 2. Die Bestellung wurde unterzeichnet und es wurde 7 Tage lang keine Bewertung durchgeführt. Wenn die Bestellung abläuft und nicht ausgewertet wird, wird die Bestellung standardmäßig positiv bewertet. Wenn der Händler die Bestellung innerhalb von 5 Minuten nicht erhält, wird die Bestellung abgebrochen Es wird eine SMS-Erinnerung gesendet ... Für Szenarien mit langen Verzögerungen und geringer Echtzeitleistung können wir die Aufgabenplanung verwenden, um eine regelmäßige Abfrageverarbeitung durchzuführen. Zum Beispiel: xxl-job Heute werden wir auswählen

So implementieren Sie verteilte Sperren mit Redis in SpringBoot So implementieren Sie verteilte Sperren mit Redis in SpringBoot Jun 03, 2023 am 08:16 AM

1. Redis implementiert das Prinzip der verteilten Sperren und warum verteilte Sperren erforderlich sind. Bevor über verteilte Sperren gesprochen wird, muss erläutert werden, warum verteilte Sperren erforderlich sind. Das Gegenteil von verteilten Sperren sind eigenständige Sperren. Wenn wir Multithread-Programme schreiben, vermeiden wir Datenprobleme, die durch den gleichzeitigen Betrieb einer gemeinsam genutzten Variablen verursacht werden. Normalerweise verwenden wir eine Sperre, um die Richtigkeit der gemeinsam genutzten Variablen sicherzustellen Die gemeinsam genutzten Variablen liegen im gleichen Prozess. Wenn es mehrere Prozesse gibt, die gleichzeitig eine gemeinsam genutzte Ressource betreiben müssen, wie können sie sich dann gegenseitig ausschließen? Heutige Geschäftsanwendungen sind in der Regel Microservice-Architekturen, was auch bedeutet, dass eine Anwendung mehrere Prozesse bereitstellen muss. Wenn mehrere Prozesse dieselbe Datensatzzeile in MySQL ändern müssen, ist eine Verteilung erforderlich, um fehlerhafte Daten zu vermeiden wird zu diesem Zeitpunkt eingeführt. Der Stil ist gesperrt. Punkte erreichen wollen

So lösen Sie das Problem, dass Springboot nach dem Einlesen in ein JAR-Paket nicht auf die Datei zugreifen kann So lösen Sie das Problem, dass Springboot nach dem Einlesen in ein JAR-Paket nicht auf die Datei zugreifen kann Jun 03, 2023 pm 04:38 PM

Springboot liest die Datei, kann aber nach dem Packen in ein JAR-Paket nicht auf die neueste Entwicklung zugreifen. Es gibt eine Situation, in der Springboot die Datei nach dem Packen in ein JAR-Paket nicht lesen kann ist ungültig und kann nur über den Stream gelesen werden. Die Datei befindet sich unter resources publicvoidtest(){Listnames=newArrayList();InputStreamReaderread=null;try{ClassPathResourceresource=newClassPathResource("name.txt");Input

So implementieren Sie Springboot+Mybatis-plus, ohne SQL-Anweisungen zum Hinzufügen mehrerer Tabellen zu verwenden So implementieren Sie Springboot+Mybatis-plus, ohne SQL-Anweisungen zum Hinzufügen mehrerer Tabellen zu verwenden Jun 02, 2023 am 11:07 AM

Wenn Springboot + Mybatis-plus keine SQL-Anweisungen zum Hinzufügen mehrerer Tabellen verwendet, werden die Probleme, auf die ich gestoßen bin, durch die Simulation des Denkens in der Testumgebung zerlegt: Erstellen Sie ein BrandDTO-Objekt mit Parametern, um die Übergabe von Parametern an den Hintergrund zu simulieren dass es äußerst schwierig ist, Multi-Table-Operationen in Mybatis-plus durchzuführen. Wenn Sie keine Tools wie Mybatis-plus-join verwenden, können Sie nur die entsprechende Mapper.xml-Datei konfigurieren und die stinkende und lange ResultMap konfigurieren Schreiben Sie die entsprechende SQL-Anweisung. Obwohl diese Methode umständlich erscheint, ist sie äußerst flexibel und ermöglicht es uns

Vergleich und Differenzanalyse zwischen SpringBoot und SpringMVC Vergleich und Differenzanalyse zwischen SpringBoot und SpringMVC Dec 29, 2023 am 11:02 AM

SpringBoot und SpringMVC sind beide häufig verwendete Frameworks in der Java-Entwicklung, es gibt jedoch einige offensichtliche Unterschiede zwischen ihnen. In diesem Artikel werden die Funktionen und Verwendungsmöglichkeiten dieser beiden Frameworks untersucht und ihre Unterschiede verglichen. Lassen Sie uns zunächst etwas über SpringBoot lernen. SpringBoot wurde vom Pivotal-Team entwickelt, um die Erstellung und Bereitstellung von Anwendungen auf Basis des Spring-Frameworks zu vereinfachen. Es bietet eine schnelle und einfache Möglichkeit, eigenständige, ausführbare Dateien zu erstellen

Wie SpringBoot Redis anpasst, um die Cache-Serialisierung zu implementieren Wie SpringBoot Redis anpasst, um die Cache-Serialisierung zu implementieren Jun 03, 2023 am 11:32 AM

1. Passen Sie den RedisTemplate1.1-Standard-Serialisierungsmechanismus an. Die API-basierte Redis-Cache-Implementierung verwendet die RedisTemplate-Vorlage für Daten-Caching-Vorgänge. Öffnen Sie hier die RedisTemplate-Klasse und zeigen Sie die Quellcodeinformationen der Klasse publicclassRedisTemplateextendsRedisAccessorimplementsRedisOperations an. Schlüssel deklarieren, verschiedene Serialisierungsmethoden des Werts, der Anfangswert ist leer @NullableprivateRedisSe

So erhalten Sie den Wert in application.yml in Springboot So erhalten Sie den Wert in application.yml in Springboot Jun 03, 2023 pm 06:43 PM

In Projekten werden häufig einige Konfigurationsinformationen benötigt. Diese Informationen können in der Testumgebung und in der Produktionsumgebung unterschiedliche Konfigurationen haben und müssen möglicherweise später basierend auf den tatsächlichen Geschäftsbedingungen geändert werden. Wir können diese Konfigurationen nicht fest im Code codieren. Am besten schreiben Sie sie in die Konfigurationsdatei. Sie können diese Informationen beispielsweise in die Datei application.yml schreiben. Wie erhält oder verwendet man diese Adresse im Code? Es gibt 2 Methoden. Methode 1: Wir können den Wert, der dem Schlüssel in der Konfigurationsdatei (application.yml) entspricht, über den mit @Value versehenen Wert erhalten. Diese Methode eignet sich für Situationen, in denen es relativ wenige Mikrodienste gibt: Tatsächlich Projekte, wenn das Geschäft kompliziert ist, Logik

See all articles