使用Spring Data Elasticsearch的PIT API时如何设置切片?
在Spring Data Elasticsearch 5.1.4中结合PIT实现切片检索
我正在开发一个Spring Boot应用,需要检索10万+文档,当前使用PIT(时间点)API来实现该需求。此前通过Elasticsearch原生库实现了带切片的PIT功能,代码如下:
SearchRequest searchRequest = new SearchRequest(); final PointInTimeBuilder pointInTimeBuilder = new PointInTimeBuilder(pitId); SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder(); searchSourceBuilder.query(query); pointInTimeBuilder.setKeepAlive(minutes); searchSourceBuilder.pointInTimeBuilder(pointInTimeBuilder); SliceBuilder sliceBuilder = new SliceBuilder(i, totalSlice); searchSourceBuilder.slice(sliceBuilder); searchRequest.source(searchSourceBuilder); RequestOptions.Builder options = RequestOptions.DEFAULT.toBuilder(); RequestConfig requestConfig = RequestConfig.custom().setContentCompressionEnabled(true).build(); options.setRequestConfig(requestConfig); SearchResponse searchResponse = this.restHighLevelClient.search(searchRequest, options.build());
目前已切换为仅使用spring-data-elasticsearch:5.1.4,且已通过Spring Data Elasticsearch官方的时间点(Point-in-Time)文档实现了PIT API,现在需要在此基础上实现切片功能。以下是对应实现方案:
实现步骤与代码示例
- 准备PIT ID:确保已通过
ElasticsearchOperations的createPointInTime方法生成有效的PIT ID(首次创建),或复用已存在的PIT ID。 - 构建带切片的PIT查询:使用
NativeQueryBuilder同时配置PIT参数与切片参数:
import org.springframework.data.elasticsearch.core.ElasticsearchOperations; import org.springframework.data.elasticsearch.core.SearchHit; import org.springframework.data.elasticsearch.core.query.NativeQuery; import org.springframework.data.elasticsearch.core.query.NativeQueryBuilder; import org.springframework.data.elasticsearch.core.query.PointInTime; import org.springframework.data.elasticsearch.core.query.Slice; // 注入ElasticsearchOperations实例 private final ElasticsearchOperations elasticsearchOperations; /** * 执行单一切片的PIT查询 * @param pitId 已创建的PIT ID * @param currentSlice 当前切片索引(从0开始) * @param totalSlices 总切片数 * @param query 查询条件 * @param keepAlive PIT存活时间(如"5m") */ public void searchWithPitAndSlice(String pitId, int currentSlice, int totalSlices, Query query, String keepAlive) { // 构建切片配置,对应原生SliceBuilder(currentSlice, totalSlices) Slice slice = Slice.of(currentSlice, totalSlices); // 组装带PIT和切片的查询 NativeQuery nativeQuery = new NativeQueryBuilder() .withQuery(query) .withPointInTime(PointInTime.of(pitId, keepAlive)) // 设置PIT ID与存活时间 .withSlice(slice) // 添加切片配置 .build(); // 执行查询并获取结果 List<SearchHit<YourDocumentEntity>> hits = elasticsearchOperations.search(nativeQuery, YourDocumentEntity.class).getSearchHits(); // 处理当前切片的检索结果... }
- 循环处理所有切片:遍历所有切片索引,完成全量数据检索后删除PIT释放资源:
/** * 处理所有切片的PIT查询 * @param pitId 已创建的PIT ID * @param query 查询条件 * @param keepAlive PIT存活时间 * @param totalSlices 总切片数 */ public void processAllSlices(String pitId, Query query, String keepAlive, int totalSlices) { for (int i = 0; i < totalSlices; i++) { searchWithPitAndSlice(pitId, i, totalSlices, query, keepAlive); } // 所有切片处理完成后,删除PIT elasticsearchOperations.deletePointInTime(pitId); }
关键说明
- 切片索引从0开始计数,总切片数建议设置为Elasticsearch集群的分片数,可最大化利用并行检索能力。
- 所有切片查询必须使用同一个PIT ID,确保基于同一时间点的快照数据检索。
- 务必在检索完成后删除PIT,避免占用集群资源。
相关参考文档:
- Spring Data Elasticsearch 官方文档:时间点(Point-in-Time)
- Elasticsearch 官方文档:时间点(Point-in-Time)API(8.3版本)
内容的提问来源于stack exchange,提问作者Omar
相关产品推荐
相关产品推荐

