Elasticsearch 7.17 PIT切片查询结果数不匹配及技术疑问
Elasticsearch 7.17 PIT+切片批量查询问题解答
场景说明
使用Elasticsearch 7.17,需高效获取约5亿条文档的ID、版本号、完整索引名(通过别名搜索)及指定字段,采用Point In Time(PIT)结合切片方案,但查询返回的命中数与总计数不匹配,相关代码如下:
启动处理方法
public void start() { ElasticPointInTime elasticPointInTime = new ElasticPointInTime(client); String indexName = "myindex"; long keepAliveMins = 1; String pointInTimeId = elasticPointInTime.getPointInTimeId(indexName, keepAliveMins); double count = elasticPointInTime.countPointInTime(pointInTimeId, keepAliveMins); if (count > 0) { int maxSlices = 10; //Don't now how to work out what this value should be Object[] searchAfterValue = null; for (int currentSlice = 0; currentSlice < maxSlices; currentSlice++) { DocumentDataResponse response = elasticPointInTime.searchPointInTime(pointInTimeId, keepAliveMins, currentSlice, maxSlices, searchAfterValue); searchAfterValue = response.getSearchAfterValue(); processDocumentData(response.getPointInTimeData()); } } }
PIT查询处理类
public class ElasticPointInTime { private static final String FIELD_NAME = "field"; private static final int MAX_SIZE = 10000; private final RestHighLevelClient client; public ElasticPointInTime(RestHighLevelClient client) { this.client = client; } public String getPointInTimeId(String indexName, long keepAliveMins) { String pointInTimeId = null; OpenPointInTimeRequest openRequest = new OpenPointInTimeRequest(indexName); openRequest.keepAlive(TimeValue.timeValueMinutes(keepAliveMins)); try { OpenPointInTimeResponse openResponse = client.openPointInTime(openRequest, RequestOptions.DEFAULT); if (openResponse != null) { pointInTimeId = openResponse.getPointInTimeId(); } } catch (IOException e) { e.printStackTrace(); } return pointInTimeId; } public long countPointInTime(String pointInTimeId, long keepAliveMins) { long count = -1; CountRequest countRequest = new CountRequest(); final PointInTimeBuilder pointInTimeBuilder = new PointInTimeBuilder(pointInTimeId); pointInTimeBuilder.setKeepAlive(TimeValue.timeValueMinutes(keepAliveMins)); SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder(); countRequest.source(searchSourceBuilder.pointInTimeBuilder(pointInTimeBuilder)); try { CountResponse countResponse = client.count(countRequest, RequestOptions.DEFAULT); if (countResponse != null) { count = countResponse.getCount(); } } catch (IOException e) { e.printStackTrace(); } return count; } public DocumentDataResponse searchPointInTime(String pointInTimeId, long keepAliveMins, int currentSlice, int maxSlices, Object[] searchAfterValue) { DocumentDataResponse response = new DocumentDataResponse(); Map<String, DocumentData> pointInTimeData = new HashMap<String, DocumentData>(); SearchRequest searchRequest = new SearchRequest(); final PointInTimeBuilder pointInTimeBuilder = new PointInTimeBuilder(pointInTimeId); pointInTimeBuilder.setKeepAlive(TimeValue.timeValueMinutes(keepAliveMins)); SliceBuilder sliceBuilder = new SliceBuilder(currentSlice, maxSlices); //Tried various fieldSorters but really just want data in order it was ingested //FieldSortBuilder sharedDocSortBuilder = SortBuilders.fieldSort("_shard_doc").order(SortOrder.ASC); FieldSortBuilder pitSortBuilder = SortBuilders.pitTiebreaker(); //FieldSortBuilder docSortBuilder = SortBuilders.fieldSort("_doc").order(SortOrder.ASC); SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder() .fetchSource(false) // Dont need all document data .version(true) // Get document version .fetchField(FIELD_NAME) // Get specific field .size(MAX_SIZE) // Set number of rows to return otherwise only 10 is returned .sort(pitSortBuilder) // Sort .pointInTimeBuilder(pointInTimeBuilder) // Point in time builder .slice(sliceBuilder) // Sets the slice builder ; if (searchAfterValue != null && searchAfterValue.length > 0) { searchSourceBuilder.searchAfter(searchAfterValue); } searchRequest.source(searchSourceBuilder); Object[] lastSearchAfterValue = null; try { SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT); if (searchResponse != null) { SearchHits hits = searchResponse.getHits(); SearchHit[] searchHits = hits.getHits(); for (SearchHit hit: searchHits) { String docId = hit.getId(); String indexName = hit.getIndex(); long versionNumber = hit.getVersion(); DocumentField docField = hit.field(FIELD_NAME); String fieldValue = (docField == null || docField.getValue() == null) ? "" : docField.getValue().toString(); DocumentData docData = new DocumentData(); docData.setDocumentId(docId); docData.setIndexName(indexName); docData.setVersionNumber(versionNumber); docData.setFieldData(fieldValue); response.getPointInTimeData().put(docId, docData); lastSearchAfterValue = hit.getSortValues(); } response.searchAfterValue = lastSearchAfterValue; } } catch (IOException e) { e.printStackTrace(); } return response; } }
问题解答
1. 结果数与总计数不匹配的原因
- 切片未循环分页:当前代码每个切片仅调用一次
searchPointInTime,仅获取一页数据(MAX_SIZE=10000),但每个切片对应的数据集远大于单页容量,需循环调用直到searchAfterValue为null(无更多数据)。 - HashMap去重问题:用
docId作为HashMap的key存储结果,若不同索引存在相同ID的文档,会被覆盖,导致计数减少。建议改用indexName + "_" + docId作为唯一key,避免覆盖。 - PIT过期风险:若
keepAliveMins设置过短,在批量查询过程中PIT可能提前过期,导致后续查询无法获取完整数据。建议根据数据量预估查询耗时,适当延长PIT的存活时间。
2. 按 ingestion 顺序返回是否需要排序?性能影响如何?
必须设置排序,且推荐使用_doc排序:
_doc排序是Elasticsearch中性能最优的排序方式,它直接遵循Lucene内部的文档ID顺序(即文档写入顺序),无需额外计算排序值,几乎无性能开销。searchAfter依赖稳定的排序规则,若不设置排序,每次请求的结果顺序可能不稳定,导致searchAfter无法正确分页,出现漏数据或重复数据的情况。- 替换当前的
pitSortBuilder为SortBuilders.fieldSort("_doc").order(SortOrder.ASC)即可满足按 ingestion 顺序返回的需求。
3. 为何必须设置size参数?
切片的作用是将整个数据集划分为多个独立的子集,不控制单页返回的数据量:
- Elasticsearch默认的单页返回条数是10,若不设置
size,每次请求仅返回10条数据,效率极低,且会导致需要大量请求才能完成数据获取。 - 设置
size=10000是合理的,这是Elasticsearch默认的单页最大返回条数(可通过修改index.max_result_window调整,但不建议过大,避免内存压力)。
4. 如何计算合适的切片数量?
切片数量建议参考以下规则:
- 优先匹配主分片数:切片数量最好等于目标索引的主分片数,这样每个切片可以对应一个主分片,最大化并行度,避免分片资源竞争。
- 不超过主分片数的2倍:若集群资源充足,可设置为分片数的2倍,但过多切片会导致额外的请求开销和资源占用,反而降低效率。
- 结合集群负载调整:若集群CPU、内存资源紧张,可适当减少切片数量(比如为主分片数的1/2),避免集群过载。
- 例如:若目标索引有8个主分片,可设置
maxSlices=8或16,根据实际测试的性能和负载情况调整。
内容的提问来源于stack exchange,提问作者karen
相关产品推荐
相关产品推荐

