Elasticsearch8.x深度分页查询
·
import co.elastic.clients.elasticsearch._types.FieldValue;
import co.elastic.clients.elasticsearch._types.SortOptions;
import co.elastic.clients.elasticsearch._types.query_dsl.Query;
import co.elastic.clients.elasticsearch._types.query_dsl.RangeQuery;
import co.elastic.clients.elasticsearch.core.SearchRequest;
import co.elastic.clients.elasticsearch.core.search.Hit;
import co.elastic.clients.json.JsonData;
import com.gy.hsxt.elasticsearch.constant.SearchResultConstant;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch.core.SearchResponse;
import co.elastic.clients.util.ObjectBuilder;
import org.springframework.stereotype.Service;
import java.math.BigDecimal;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
/**
* @author Administrator
* @Description:
* @date 2025/10/13 9:26
*/
@Slf4j
@Service
public class ElasticsearchGenericSearchService {
// 注入ES客户端(通用服务无需绑定具体索引,客户端全局复用)
@Autowired
private ElasticsearchClient esClient;
// 分页最大值
private static final int PAGE_SIZE_MAX= 1000;
private static final int PAGE_SIZE_MIN= BigDecimal.ONE.intValue();
public <T> Map<String,Object> pageMapSearch(
String indexName,
int pageNo,
int pageSize,
List<SortOptions> sortOptionsList,
Function<Query.Builder, ObjectBuilder<Query>> queryBuilder,
Class<T> dataClass
) throws Exception {
Map<String,Object> dataMap = new HashMap<>();
// 1. 参数校验(通用逻辑,所有类型复用)
validateParams(indexName, pageNo, pageSize, dataClass);
// 2. 分页核心逻辑(Search After通用实现,与具体数据类型解耦)
int targetPage = pageNo;
int currentPage = 1;
long total = 0;
List<T> targetData = new ArrayList<>();
AtomicReference<List<FieldValue>> searchAfterRef = new AtomicReference<>(null);
boolean hasMore = true;
while (hasMore && currentPage <= targetPage) {
// 3. 构建通用查询请求(索引、排序、查询条件均由外部传入)
SearchRequest request = buildGenericSearchRequest(
indexName, pageSize, searchAfterRef.get(), sortOptionsList, queryBuilder
);
// 4. 执行查询(泛型适配不同数据类型)
SearchResponse<T> response = esClient.search(request, dataClass);
log.info("查询结果:{}条 sql:{}",response.hits().total().value(),request.toString());
// 5. 提取数据(通用逻辑)
List<T> currentData = extractData(response);
// 6. 目标页判断(通用逻辑)
if (currentPage == targetPage) {
targetData = currentData;
total = response.hits().total().value();
}
// 7. 更新分页状态(通用逻辑)
hasMore = updatePaginationState(response, pageSize, searchAfterRef);
currentPage++;
}
dataMap.put(SearchResultConstant.DATA, targetData);
// 获取总数
dataMap.put(SearchResultConstant.TOTAL, total);
return dataMap;
}
/**
* 通用分页查询(Search After模式,支持任意索引、任意数据类型)
* @param indexName 目标索引名(如"order_index"/"product_index")
* @param pageNo 目标页码(从1开始)
* @param pageSize 每页条数
* @param sortOptionsList 排序条件构造器(外部传入,适配不同数据类型的排序字段)
* @param queryBuilder 查询条件构造器(外部传入,适配不同业务的过滤逻辑)
* @param dataClass 数据类型Class(如OrderIndex.class/ProductIndex.class)
* @return 目标页数据列表(非null)
* @param <T> 数据类型泛型
* @throws Exception 查询异常(上层处理)
*/
public <T> List<T> pageSearch(
String indexName,
int pageNo,
int pageSize,
List<SortOptions> sortOptionsList,
Function<Query.Builder, ObjectBuilder<Query>> queryBuilder,
Class<T> dataClass
) throws Exception {
Map<String,Object> dataMap = pageMapSearch(
indexName,
pageNo,
pageSize,
sortOptionsList,
queryBuilder,
dataClass
);
return (List<T>) dataMap.get(SearchResultConstant.DATA);
}
// ------------------------------ 通用工具方法(参数校验、请求构建等)------------------------------
private <T> void validateParams(String indexName, int pageNo, int pageSize, Class<T> dataClass) {
if (indexName == null || indexName.trim().isEmpty()) {
throw new IllegalArgumentException("索引名indexName不能为空");
}
if (pageNo < PAGE_SIZE_MIN) {
throw new IllegalArgumentException("页码pageNo必须≥1,当前值:" + pageNo);
}
if (pageSize < PAGE_SIZE_MIN || pageSize > PAGE_SIZE_MAX) {
throw new IllegalArgumentException("页大小pageSize必须在1-1000之间,当前值:" + pageSize);
}
if (dataClass == null) {
throw new IllegalArgumentException("数据类型Class不能为空");
}
}
private <T> SearchRequest buildGenericSearchRequest(
String indexName,
int pageSize,
List<FieldValue> searchAfter,
List<SortOptions> sortOptionsList,
Function<Query.Builder, ObjectBuilder<Query>> queryBuilder
) {
return SearchRequest.of(s -> {
s.index(indexName)
.query(queryBuilder) // 外部传入查询条件(如订单查商户ID、商品查分类)
.size(pageSize)
.from(0);
if(sortOptionsList != null && !sortOptionsList.isEmpty()){
s.sort(sortOptionsList); // 外部传入排序条件
}
if (searchAfter != null && !searchAfter.isEmpty()) {
s.searchAfter(searchAfter);
}
return s;
});
}
private <T> List<T> extractData(SearchResponse<T> response) {
if (response.hits() == null || response.hits().hits() == null) {
return new ArrayList<>();
}
return response.hits().hits().stream()
.map(hit -> hit.source())
.filter(java.util.Objects::nonNull)
.toList();
}
private <T> boolean updatePaginationState(
SearchResponse<T> response,
int pageSize,
AtomicReference<List<FieldValue>> searchAfterRef
) {
List<Hit<T>> hits = response.hits().hits();
if (hits == null || hits.size() < pageSize) {
return false; // 数据不足,无下一页
}
// 更新searchAfter为最后一条数据的排序值
List<FieldValue> lastSort = hits.get(hits.size() - 1).sort();
searchAfterRef.set(lastSort);
return true;
}
/**
* 添加时间范围查询条件到布尔查询构建器中
* 支持开始时间和结束时间的范围查询
*
* @param fieldName 查询参数字段名
* @param startTime 查询开始时间
* @param endTime 查询结束时间
* @see RangeQuery
*/
public RangeQuery addTimeRangeQuery(String fieldName, Long startTime, Long endTime) {
// 时间范围查询优化版
if (startTime != null || endTime != null) {
RangeQuery.Builder rangeBuilder = new RangeQuery.Builder();
// 设置字段名
rangeBuilder.field(fieldName);
// 添加开始时间条件(大于等于)
if (startTime != null) {
rangeBuilder.gte(JsonData.of(startTime));
}
// 添加结束时间条件(小于等于)
if (endTime != null) {
rangeBuilder.lte(JsonData.of(endTime));
}
return rangeBuilder.build();
}
return null;
}
}
调用示例:
@Slf4j
@Service
public class SearchOrderServiceImpl {
/**
* 索引名称常量
*/
private static final String INDEX_NAME = "order_index";
@Autowired
private ElasticsearchGenericSearchService elasticsearchGenericSearchService;
public List<SearchOrderRespVO> pageQueryOrderListByEntCustId(AbstractOrderQuery params) {
if(params.getCurPage() <= BigDecimal.ZERO.intValue())
params.setCurPage(Constants.DEFAULT_PAGE_NUM);
// 1. 定义排序字段(必须与 searchAfter 的值列表对应)
List<SortOptions> sortOptions = new ArrayList<>();
// 排序字段1:create_time(长整型)
// 排序条件:先按 createDate 升序,再按 id.keyword 升序(订单专属排序)
sortOptions.add(SortOptions.of(s -> s
.field(f -> f.field("createDate").order(SortOrder.Desc))
));
// 排序字段2:_id(字符串,确保唯一性)
sortOptions.add(SortOptions.of(s -> s
.field(f -> f.field("id.keyword").order(SortOrder.Desc))
));
// 查询条件:商户ID=MERCHANT_001(订单专属过滤)
Function<Query.Builder, ObjectBuilder<Query>> orderQuery = q ->{
q.bool(b -> b
.must(m -> m.term(t -> t.field("userId").value(params.getUserId())))
.must(m -> m.term(t -> t.field("mobile").value(params.getMobile())))
.must(m -> m.range(elasticsearchGenericSearchService.addTimeRangeQuery("createDate", params.getStartTime(), params.getEndTime())))
);
return q;
};
List<SearchOrderRespVO> vos = new ArrayList<>();
try {
List<OrderIndex> list = elasticsearchGenericSearchService.pageSearch(
INDEX_NAME, params.getCurPage(), 10, // 第1页,每页10条
sortOptions, orderQuery,
OrderIndex.class // 订单数据类型
);
// 遍历结果
list.forEach(order -> {
System.out.println("count======="+order );
});
}catch (ElasticsearchException e) {
log.error("调用ES查询订单列表失败,参数:{}",params, e);
} catch (Exception e) {
log.error("查询订单列表失败,参数:{}",params, e);
}
return vos;
}
}
pom.xml 依赖配置(JDK21)
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-elasticsearch</artifactId>
<version>${spring-boot.version}</version>
</dependency>
<!-- Kafka -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.elasticsearch</groupId>
<artifactId>elasticsearch</artifactId>
<version>8.15.1</version>
</dependency>
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<version>7.17.26</version>
</dependency>
更多推荐


所有评论(0)