为了兼顾 Scroll 的数据一致性视图 与 search_after 的轻量低内存消耗,Elasticsearch 自 7.10 起正式推出了 Point in Time (PIT,时间点视图)。

一、为什么纯 search_after 还会出现“漏数”或“重数”?
search_after 默认是强实时的。这意味着在分页连续滚动的过程中,底层索引可能正在经历并发写操作:

时间点 T1: 客户端读取 Page 1 (末尾游标: create_time = 10:00:05, id = 88)
│
时间点 T2: 业务高频并发写入:插入一条历史补单数据 (create_time = 10:00:03, id = 99)
更新了一条正在排序区间的数据,使其排序值发生改变
触发了一次 Lucene Segment Merge(段合并与已删除文档清理)
│
时间点 T3: 客户端带上 (10:00:05, 88) 游标请求 Page 2
此时新插入的 id = 99 永远不会被遍历到,发生【静默漏数】!
痛点根源:纯 search_after 没有固定查询基准时间,翻页跨度越长,前后的视图越脱节。

Scroll 的代价:传统 Scroll 会为整个查询生成一个重型 Context,强行保留旧的 Segment 不被合并删除,导致磁盘、文件句柄与内存急剧消耗。

PIT 的诞生:PIT 提供了一个轻量级的索引时间点视图。它锁定了打开时间点时刻的分片与 Lucene 段引用,确保后续所有关联该 PIT 的滚动检索看到的是一份完全冻结一致的数据视图,不受后续写入、删除或段合并的影响。

二、PIT 核心底层原理:它比 Scroll 轻在哪里?
[客户端] ──(1. open PIT)──► [协调节点]
│ (广播至相关分片)
▼
[Lucene Segment 引用计数 (Ref Count)]
┌─────────────────────┴─────────────────────┐
▼ ▼
[Segment 1 (活动)] [Segment 2 (活动)]
Ref Count: 1 -> 2 (保留引用) Ref Count: 1 -> 2 (保留引用)
│ │
└──────────────┬────────────────────────────┘
▼
后续即使后台触发 Segment Merge 合成 Segment 3,旧段文件在磁盘上仅标记
逻辑删除,只要 PIT 的 Ref Count > 0,操作系统就不会物理解绑删除旧文件!
轻量级段引用(Segment Reader Pinning):
PIT 本质上只是为分片底层的 Lucene IndexReader 增加了一个引用计数(Keep-Alive)。它并不像 Scroll 那样需要常驻整套完整的查询上下文(Search Context)和协调节点内存聚合状态。

解耦查询条件与视图状态:
在 Scroll 中,查询语句和排序规则必须在初始化时定死;而 PIT 是一个独立的视图载体。你可以在同一个 PIT 视图内,使用不同的 Query、不同的 Filter 和不同的排序进行多次灵活检索。

支持跨索引路由持久化:
PIT 生成的字符串 ID 经过 Base64 编码,内部包含了分片物理位置和段版本。协调节点拿到 PIT 后,能精确路由到当时负责对应数据分片的节点。

三、RESTful 核心调用规范

  1. 打开一个 PIT 视图
    在目标索引上创建时间点,并指定存活时间(keep_alive):

HTTP

POST /idx_orders/_pit?keep_alive=2m
响应返回一个唯一的 id:

JSON

{
“id”: “46ToAwEKaWR4X29yZGVycxZlR…”,
“_shards”: {
“total”: 3,
“successful”: 3,
“failed”: 0
}
}
2. 基于 PIT + search_after 分页滚动
此时检索请求不再在 URL 中指定索引名称,而是在请求体(Body)内注入 pit 参数:

HTTP

POST /_search
{
“size”: 10,
“pit”: {
“id”: “46ToAwEKaWR4X29yZGVycxZlR…”,
“keep_alive”: “2m”
},
“query”: {
“match”: { “status”: “PAID” }
},
“sort”: [
{ “create_time”: “desc” },
{ “_shard_doc”: “asc” } // PIT 专属的高效排序字段
],
“search_after”: [ 1773801234000, 42 ]
}
性能利器 _shard_doc:
在配合 PIT 使用时,官方强烈推荐使用隐式的 _shard_doc 作为排序兜底键(替代 _id)。_shard_doc 按照内部 Lucene 内部 DocId 与分片顺序排序,执行速度极快,且零额外内存消耗。

  1. 显式释放 PIT 视图(非常重要)
    不要等待 keep_alive 超时自动释放,批量任务完成后必须显式关闭:

HTTP

DELETE /_pit
{
“id”: “46ToAwEKaWR4X29yZGVycxZlR…”
}
四、Spring Boot 3 + Elasticsearch Java Client 工业级实现
以下基于官方主流的 Elasticsearch Java Client (8.x+) 实现支持 PIT 的安全滚动拉取组件:

  1. 深度滚动服务代码
    Java

package com.example.es.pit;

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.FieldValue;
import co.elastic.clients.elasticsearch._types.SortOptions;
import co.elastic.clients.elasticsearch._types.SortOrder;
import co.elastic.clients.elasticsearch._types.Time;
import co.elastic.clients.elasticsearch.core.*;
import co.elastic.clients.elasticsearch.core.search.Hit;
import co.elastic.clients.elasticsearch.core.search.PointInTimeReference;
import com.example.es.model.OrderDocument;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.util.CollectionUtils;

import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.function.Consumer;

@Slf4j
@Service
@RequiredArgsConstructor
public class PitScrollExportService {

private final ElasticsearchClient esClient;
private static final String INDEX_NAME = "idx_orders";
private static final Time PIT_KEEP_ALIVE = Time.of(t -> t.time("2m"));

/**
 * 流式全量导出/滚动遍历(规避内存 OOM 与并发写干扰)
 *
 * @param batchSize 每次拉取批次大小(推荐 500~2000)
 * @param consumer  批次数据处理回调逻辑
 */
public void scrollAllOrders(int batchSize, Consumer<List<OrderDocument>> consumer) {
    String pitId = null;
    try {
        // 1. 初始化打开 PIT 视图
        OpenPointInTimeResponse openPitResp = esClient.openPointInTime(b -> b
                .index(INDEX_NAME)
                .keepAlive(PIT_KEEP_ALIVE)
        );
        pitId = openPitResp.id();
        log.info("成功创建 PIT 视图, pitId: {}", pitId);

        List<FieldValue> searchAfterValues = null;
        boolean hasMore = true;

        while (hasMore) {
            final String currentPitId = pitId;
            final List<FieldValue> currentSearchAfter = searchAfterValues;

            // 2. 构建基于 PIT 的检索请求
            SearchRequest searchRequest = SearchRequest.of(s -> {
                s.size(batchSize)
                 .pit(PointInTimeReference.of(p -> p.id(currentPitId).keepAlive(PIT_KEEP_ALIVE)))
                 // 推荐使用 _shard_doc 排序,极大降低 CPU/内存消耗并保证绝对唯一性
                 .sort(Collections.singletonList(
                         SortOptions.of(so -> so.field(f -> f.field("_shard_doc").order(SortOrder.Asc)))
                 ));

                if (!CollectionUtils.isEmpty(currentSearchAfter)) {
                    s.searchAfter(currentSearchAfter);
                }
                return s;
            });

            SearchResponse<OrderDocument> response = esClient.search(searchRequest, OrderDocument.class);

            // 注意:每次 search 响应可能会返回更新后的 pit_id,必须滚动替换
            if (response.pitId() != null) {
                pitId = response.pitId();
            }

            List<Hit<OrderDocument>> hits = response.hits().hits();
            if (CollectionUtils.isEmpty(hits)) {
                hasMore = false;
                break;
            }

            // 3. 提取业务数据交由回调消费者消费
            List<OrderDocument> documents = new ArrayList<>(hits.size());
            for (Hit<OrderDocument> hit : hits) {
                if (hit.source() != null) {
                    documents.add(hit.source());
                }
            }
            consumer.accept(documents);

            // 4. 提取当前批次最后一条的 sort 值,充当下一次滚动的 searchAfter
            Hit<OrderDocument> lastHit = hits.get(hits.size() - 1);
            searchAfterValues = lastHit.sort();

            // 若拉取量小于 batchSize,说明已触达视图末尾
            if (hits.size() < batchSize) {
                hasMore = false;
            }
        }
    } catch (Exception e) {
        log.error("PIT 滚动导出发生异常", e);
        throw new RuntimeException("Export failed via PIT", e);
    } finally {
        // 5. 必须在 finally 中关闭释放 PIT,归还段句柄
        if (pitId != null) {
            closePit(pitId);
        }
    }
}

private void closePit(String pitId) {
    try {
        final String finalPitId = pitId;
        ClosePointInTimeResponse closeResp = esClient.closePointInTime(b -> b.id(Collections.singletonList(finalPitId)));
        log.info("成功关闭 PIT 视图, closedShards: {}", closeResp.numFreed());
    } catch (IOException e) {
        log.warn("关闭 PIT 异常,等待 TTL 自动驱逐, pitId: {}", pitId, e);
    }
}

}
五、生产落地必须牢记的避坑要点
响应中的 pit_id 必须动态更新:
在循环滚动中,每次 SearchResponse 会返回一个新的 pit_id(可能与初始传入的值相同或刷新了上下文)。必须用响应中最新的 response.pitId() 覆盖原变量,否则在长时间滚动中视图可能失效。

keep_alive 绝不要设得过大:
keep_alive 只需要略大于“单次批处理耗时”(例如 1m 或 2m)。因为每次发出包含该 PIT 的查询时,其存活时间都会被重新刷新(续期)。盲目设置数小时会导致旧 Segment 无法合并释放,打爆磁盘。

优先选择 _shard_doc 排序:
如果导出需求对业务排序没有严苛要求(只是为了全量拉取),使用 _shard_doc 的遍历效率远高于普通字段,它能直接绕过字段 DocValues 的跨分片归并开销。

清理的兜底措施:
务必确保将 closePointInTime 放在 Java 的 finally 块中。若发生因网络中断导致的孤立 PIT,ES 会在 keep_alive 到期后自动执行 GC 清理,无需过度惊慌。

Logo

智能硬件社区聚焦AI智能硬件技术生态,汇聚嵌入式AI、物联网硬件开发者,打造交流分享平台,同步全国赛事资讯、开展 OPC 核心人才招募,助力技术落地与开发者成长。

更多推荐