Table des matières
一、引入依赖jar
二、application.properties配置
三、使用
其它
Maison Java javaDidacticiel Comment utiliser RestHighLevelClient dans springboot2+es7

Comment utiliser RestHighLevelClient dans 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>
Copier après la connexion

二、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
Copier après la connexion

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);
                }
            }
        }
    }
}
Copier après la connexion

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;
    }
}
Copier après la connexion

三、使用

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;
    }
 
}
Copier après la connexion

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查询返回的状态码异常!");
        }
    }
}
Copier après la connexion

其中,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;
    }
 
}
Copier après la connexion
/**
 * 聚合桶对象
 *
 * @author yangzihe
 * @date 2022/1/26
 */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class AggregationBucketPO {
 
    /**
     * 聚合Bucket的key名
     */
    private String key;
 
    /**
     * 聚合Bucket的文档数量
     */
    private Long docCount;
 
    /**
     * 文档总数量
     */
    private Long docTotal;
 
}
Copier après la connexion

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());
    }
}
Copier après la connexion

其中,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;
}
Copier après la connexion
@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;
    }
}
Copier après la connexion

二级聚合分桶测试代码

    @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);
    }
Copier après la connexion

其它

如果没有用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);
    }
 
}
Copier après la connexion

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>
Copier après la connexion

聚合查询的测试用例:

@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);
    }
}
Copier après la connexion

Ce qui précède est le contenu détaillé de. pour plus d'informations, suivez d'autres articles connexes sur le site Web de PHP en chinois!

Déclaration de ce site Web
Le contenu de cet article est volontairement contribué par les internautes et les droits d'auteur appartiennent à l'auteur original. Ce site n'assume aucune responsabilité légale correspondante. Si vous trouvez un contenu suspecté de plagiat ou de contrefaçon, veuillez contacter admin@php.cn

Outils d'IA chauds

Undresser.AI Undress

Undresser.AI Undress

Application basée sur l'IA pour créer des photos de nu réalistes

AI Clothes Remover

AI Clothes Remover

Outil d'IA en ligne pour supprimer les vêtements des photos.

Undress AI Tool

Undress AI Tool

Images de déshabillage gratuites

Clothoff.io

Clothoff.io

Dissolvant de vêtements AI

AI Hentai Generator

AI Hentai Generator

Générez AI Hentai gratuitement.

Article chaud

R.E.P.O. Crystals d'énergie expliqués et ce qu'ils font (cristal jaune)
4 Il y a quelques semaines By 尊渡假赌尊渡假赌尊渡假赌
R.E.P.O. Meilleurs paramètres graphiques
4 Il y a quelques semaines By 尊渡假赌尊渡假赌尊渡假赌
R.E.P.O. Comment réparer l'audio si vous n'entendez personne
1 Il y a quelques mois By 尊渡假赌尊渡假赌尊渡假赌
R.E.P.O. Commandes de chat et comment les utiliser
1 Il y a quelques mois By 尊渡假赌尊渡假赌尊渡假赌

Outils chauds

Bloc-notes++7.3.1

Bloc-notes++7.3.1

Éditeur de code facile à utiliser et gratuit

SublimeText3 version chinoise

SublimeText3 version chinoise

Version chinoise, très simple à utiliser

Envoyer Studio 13.0.1

Envoyer Studio 13.0.1

Puissant environnement de développement intégré PHP

Dreamweaver CS6

Dreamweaver CS6

Outils de développement Web visuel

SublimeText3 version Mac

SublimeText3 version Mac

Logiciel d'édition de code au niveau de Dieu (SublimeText3)

Comment Springboot intègre Jasypt pour implémenter le chiffrement des fichiers de configuration Comment Springboot intègre Jasypt pour implémenter le chiffrement des fichiers de configuration Jun 01, 2023 am 08:55 AM

Introduction à Jasypt Jasypt est une bibliothèque Java qui permet à un développeur d'ajouter des fonctionnalités de chiffrement de base à son projet avec un minimum d'effort et ne nécessite pas une compréhension approfondie du fonctionnement du chiffrement. Haute sécurité pour le chiffrement unidirectionnel et bidirectionnel. technologie de cryptage basée sur des normes. Cryptez les mots de passe, le texte, les chiffres, les binaires... Convient pour l'intégration dans des applications basées sur Spring, API ouverte, pour une utilisation avec n'importe quel fournisseur JCE... Ajoutez la dépendance suivante : com.github.ulisesbocchiojasypt-spring-boot-starter2 1.1. Les avantages de Jasypt protègent la sécurité de notre système. Même en cas de fuite du code, la source de données peut être garantie.

Comment SpringBoot intègre Redisson pour implémenter la file d'attente différée Comment SpringBoot intègre Redisson pour implémenter la file d'attente différée May 30, 2023 pm 02:40 PM

Scénario d'utilisation 1. La commande a été passée avec succès mais le paiement n'a pas été effectué dans les 30 minutes. Le paiement a expiré et la commande a été automatiquement annulée 2. La commande a été signée et aucune évaluation n'a été effectuée pendant 7 jours après la signature. Si la commande expire et n'est pas évaluée, le système donne par défaut une note positive. 3. La commande est passée avec succès. Si le commerçant ne reçoit pas la commande pendant 5 minutes, la commande est annulée. 4. Le délai de livraison expire et. un rappel par SMS est envoyé... Pour les scénarios avec des délais longs et de faibles performances en temps réel, nous pouvons utiliser la planification des tâches pour effectuer un traitement d'interrogation régulier. Par exemple : xxl-job Aujourd'hui, nous allons choisir

Comment utiliser Redis pour implémenter des verrous distribués dans SpringBoot Comment utiliser Redis pour implémenter des verrous distribués dans SpringBoot Jun 03, 2023 am 08:16 AM

1. Redis implémente le principe du verrouillage distribué et pourquoi les verrous distribués sont nécessaires. Avant de parler de verrous distribués, il est nécessaire d'expliquer pourquoi les verrous distribués sont nécessaires. Le contraire des verrous distribués est le verrouillage autonome. Lorsque nous écrivons des programmes multithreads, nous évitons les problèmes de données causés par l'utilisation d'une variable partagée en même temps. Nous utilisons généralement un verrou pour exclure mutuellement les variables partagées afin de garantir l'exactitude de celles-ci. les variables partagées. Son champ d’utilisation est dans le même processus. S’il existe plusieurs processus qui doivent exploiter une ressource partagée en même temps, comment peuvent-ils s’exclure mutuellement ? Les applications métier d'aujourd'hui sont généralement une architecture de microservices, ce qui signifie également qu'une application déploiera plusieurs processus si plusieurs processus doivent modifier la même ligne d'enregistrements dans MySQL, afin d'éviter les données sales causées par des opérations dans le désordre, les besoins de distribution. à introduire à ce moment-là. Le style est verrouillé. Vous voulez marquer des points

Comment résoudre le problème selon lequel Springboot ne peut pas accéder au fichier après l'avoir lu dans un package jar Comment résoudre le problème selon lequel Springboot ne peut pas accéder au fichier après l'avoir lu dans un package jar Jun 03, 2023 pm 04:38 PM

Springboot lit le fichier, mais ne peut pas accéder au dernier développement après l'avoir empaqueté dans un package jar. Il existe une situation dans laquelle Springboot ne peut pas lire le fichier après l'avoir empaqueté dans un package jar. La raison en est qu'après l'empaquetage, le chemin virtuel du fichier. n’est pas valide et n’est accessible que via le flux Read. Le fichier se trouve sous les ressources publicvoidtest(){Listnames=newArrayList();InputStreamReaderread=null;try{ClassPathResourceresource=newClassPathResource("name.txt");Input

Comment implémenter Springboot+Mybatis-plus sans utiliser d'instructions SQL pour ajouter plusieurs tables Comment implémenter Springboot+Mybatis-plus sans utiliser d'instructions SQL pour ajouter plusieurs tables Jun 02, 2023 am 11:07 AM

Lorsque Springboot+Mybatis-plus n'utilise pas d'instructions SQL pour effectuer des opérations d'ajout de plusieurs tables, les problèmes que j'ai rencontrés sont décomposés en simulant la réflexion dans l'environnement de test : Créez un objet BrandDTO avec des paramètres pour simuler le passage des paramètres en arrière-plan. qu'il est extrêmement difficile d'effectuer des opérations multi-tables dans Mybatis-plus. Si vous n'utilisez pas d'outils tels que Mybatis-plus-join, vous pouvez uniquement configurer le fichier Mapper.xml correspondant et configurer le ResultMap malodorant et long, puis. écrivez l'instruction SQL correspondante Bien que cette méthode semble lourde, elle est très flexible et nous permet de

Comparaison et analyse des différences entre SpringBoot et SpringMVC Comparaison et analyse des différences entre SpringBoot et SpringMVC Dec 29, 2023 am 11:02 AM

SpringBoot et SpringMVC sont tous deux des frameworks couramment utilisés dans le développement Java, mais il existe des différences évidentes entre eux. Cet article explorera les fonctionnalités et les utilisations de ces deux frameworks et comparera leurs différences. Tout d’abord, découvrons SpringBoot. SpringBoot a été développé par l'équipe Pivotal pour simplifier la création et le déploiement d'applications basées sur le framework Spring. Il fournit un moyen rapide et léger de créer des fichiers exécutables autonomes.

Comment SpringBoot personnalise Redis pour implémenter la sérialisation du cache Comment SpringBoot personnalise Redis pour implémenter la sérialisation du cache Jun 03, 2023 am 11:32 AM

1. Personnalisez RedisTemplate1.1, mécanisme de sérialisation par défaut RedisAPI. L'implémentation du cache Redis basée sur l'API utilise le modèle RedisTemplate pour les opérations de mise en cache des données. Ici, ouvrez la classe RedisTemplate et affichez les informations sur le code source de la classe. Déclarer la clé, diverses méthodes de sérialisation de la valeur, la valeur initiale est vide @NullableprivateRedisSe

Comment obtenir la valeur dans application.yml au Springboot Comment obtenir la valeur dans application.yml au Springboot Jun 03, 2023 pm 06:43 PM

Dans les projets, certaines informations de configuration sont souvent nécessaires. Ces informations peuvent avoir des configurations différentes dans l'environnement de test et dans l'environnement de production, et peuvent devoir être modifiées ultérieurement en fonction des conditions commerciales réelles. Nous ne pouvons pas coder en dur ces configurations dans le code. Il est préférable de les écrire dans le fichier de configuration. Par exemple, vous pouvez écrire ces informations dans le fichier application.yml. Alors, comment obtenir ou utiliser cette adresse dans le code ? Il existe 2 méthodes. Méthode 1 : Nous pouvons obtenir la valeur correspondant à la clé dans le fichier de configuration (application.yml) via le ${key} annoté avec @Value. Cette méthode convient aux situations où il y a relativement peu de microservices. Méthode 2 : En réalité. projets, Quand les affaires sont compliquées, la logique

See all articles