数据库基础体系 · 第 127/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。
Elasticsearch 深分页与一致性:search_after、PIT、Scroll 和导出
在 Elasticsearch 中,“分页”并不只有一种含义:
- 面向用户的列表翻页,通常希望结果稳定、延迟可控;
- 批处理需要遍历大量文档,重点是吞吐和断点恢复;
- 数据导出还要求明确一致性边界,避免漏数据、重复数据或把不同时间点的数据混在一起。
search_after、PIT、Scroll 分别解决了不同层面的问题:
search_after:按照上一页最后一条文档继续查找;- PIT(Point in Time):固定一次查询所看到的索引视图;
- Scroll:保存服务端搜索上下文,适合连续批量读取;
- 导出:不是一个单独的 API,而是对查询一致性、游标生命周期、失败恢复和输出幂等性的综合设计。
这几个概念经常一起出现,但不能互相替代。search_after 本身不保证跨请求的一致视图,PIT 本身也不负责输出文件的事务提交,Scroll 也不是数据库事务。
一、先定义“深分页”和“查询一致性”
1. from + size 的普通分页
最直观的分页方式是:
GET /products/_search
{
"from": 1000,
"size": 20,
"query": {
"match": {
"name": "keyboard"
}
},
"sort": [
{
"updated_at": "desc"
}
]
}
这表示跳过前 1000 条,返回接下来的 20 条。
在分片索引中,查询不是由某一个节点直接完成的。假设索引有 3 个主分片:
- 协调节点把查询发送到 3 个分片;
- 每个分片都需要找出本地排序靠前的
from + size = 1020条候选; - 协调节点合并各分片候选;
- 再跳过前 1000 条,返回 20 条。
因此,页码越深,每个分片需要保留和排序的候选越多。粗略地说,协调过程的候选规模与下面的量相关:
其中:
- :分片数量;
- :
from; - :
size。
这不是精确的内存复杂度公式,因为实际执行还取决于查询类型、排序字段、缓存和 Lucene 实现,但它准确表达了一个工程事实:from 越大,深分页成本越高。
Elasticsearch 默认限制 from + size 不超过 index.max_result_window,通常是 10000。这个限制不是说第 10001 条数据不存在,而是防止普通分页无界地消耗协调节点和分片资源。
提高 index.max_result_window 可以绕过这个保护,但不会改变算法成本,也不会把深分页变成高效操作。
2. “一致性”至少包含三层含义
讨论深分页时,应该区分以下三种一致性。
2.1 单次请求内的一致性
一次 _search 请求会在参与查询的分片上完成一次搜索。这个范围内,协调节点会合并各分片返回的结果。
2.2 多次分页请求之间的一致性
第一页和第二页是否基于同一个索引视图?
如果两次请求之间发生了 refresh、写入、更新或删除,第二页可能看到和第一页不同的文档集合。结果可能出现:
- 同一文档在不同页重复出现;
- 某些文档被跳过;
- 文档排序位置发生变化;
- 一次分页过程中混入不同时间点的数据。
2.3 输出过程的一致性
即使 Elasticsearch 返回的每一页都没有遗漏,写入文件或下游系统时仍然可能失败:
- 第 10 页已经写入文件;
- 程序在保存进度前崩溃;
- 重启后从第 10 页重新读取;
- 文件中出现重复数据。
这已经不是 Elasticsearch 查询一致性,而是导出程序与输出介质之间的提交问题。
二、为什么 from + size 不适合深分页
Elasticsearch 的搜索结果来自 Lucene 分片。排序字段不是全局排序,而是每个分片先局部排序,再由协调节点合并成全局结果。
例如有两个分片,按照 score desc 排序:
分片 A:100, 90, 80, 70
分片 B:95, 85, 75, 65
请求 from = 2, size = 2 时,全局结果是:
100, 95, 90, 85, 80, 75, 70, 65
协调节点必须先获得足够多的分片候选,才能判断全局第 3、4 条是什么。若 from = 1_000_000,每个分片仍然不能只返回第 1_000_001 条,因为它无法预先知道其他分片会提供哪些更靠前的文档。
这解释了两件事:
- 深分页的成本主要发生在“找出前面那些不返回给用户的结果”;
- 仅仅把
size设得很小,并不能消除from带来的成本。
此外,排序还需要稳定的次序。如果只按一个非唯一字段排序,例如:
"sort": [
{
"updated_at": "desc"
}
]
许多文档可能具有相同的 updated_at。这些文档之间的顺序可能由 Lucene 内部文档编号等因素决定,而内部编号不应被当作业务稳定顺序。
三、search_after:用游标代替页码
1. 基本语义
search_after 不表示“跳到第 N 页”,而表示:
根据上一页最后一条命中的排序值,只返回排序位置在它之后的文档。
第一页:
GET /products/_search
{
"size": 2,
"query": {
"match_all": {}
},
"sort": [
{
"updated_at": "desc"
},
{
"product_id": "asc"
}
]
}
假设返回:
{
"hits": {
"hits": [
{
"_id": "p100",
"sort": ["2025-01-10T12:00:00.000Z", "p100"]
},
{
"_id": "p101",
"sort": ["2025-01-10T11:00:00.000Z", "p101"]
}
]
}
}
下一页把上一页最后一条的完整 sort 数组传入:
GET /products/_search
{
"size": 2,
"query": {
"match_all": {}
},
"sort": [
{
"updated_at": "desc"
},
{
"product_id": "asc"
}
],
"search_after": [
"2025-01-10T11:00:00.000Z",
"p101"
]
}
这里的数组必须:
- 与
sort字段一一对应; - 顺序一致;
- 使用上一页最后一条命中的实际
sort值; - 不能只传
_id,也不能传当前页第一条的排序值。
2. search_after 的排序条件
假设排序键是:
其中:
- :文档的更新时间,按降序;
- :业务唯一标识,按升序。
对于当前游标 ,下一页应返回满足下列条件的文档:
或者:
这就是多字段排序的字典序。第二个字段的作用是打破第一个字段的并列。
因此,稳定分页通常需要一个真正唯一的排序键。例如:
"sort": [
{
"updated_at": {
"order": "desc"
}
},
{
"product_id": {
"order": "asc"
}
}
]
product_id 应该是具有 keyword 类型和 doc values 的业务唯一字段,并且在分页涉及的索引范围内真正唯一。
不要直接把 _id 当作默认的高性能排序字段。_id 的索引访问语义和普通带 doc values 的 keyword 字段不同,生产设计中更适合显式建立一个用于排序和游标的业务字段。
3. search_after 的优势和限制
search_after 的主要优势是:
- 不需要计算任意深度之前的所有结果;
- 不受
from + size的深度分页方式限制; - 每一页只需要携带上一页的游标;
- 可以逐页流式处理。
但它也有明确限制:
- 不能直接跳到第 100 页;
- 不能依靠页码随机访问;
- 游标只对相同查询和排序语义有意义;
- 没有 PIT 时,跨请求不自动保证同一份数据视图;
- 如果排序键不稳定,仍可能漏数据或重复数据。
因此,search_after 解决的是“如何继续向后查找”,不是“如何固定数据版本”。
四、没有 PIT 时,search_after 为什么仍可能不一致
假设按 updated_at desc, product_id asc 分页。
初始可见数据:
A: 10, B: 9, C: 8, D: 7
第一页 size = 2 返回:
A, B
游标是 B 的排序值,即 9。
在请求下一页前,发生了两种可能的变化。
情况一:插入一个更新更晚的文档
插入:
X: 11
第二次查询使用 search_after = 9,通常仍然返回:
C, D
新文档 X 排在游标之前,因此不会出现在这次分页中。若业务希望“从分页开始那一刻起的完整结果”,它已经被漏掉。
情况二:更新了已经返回的文档
把 A 的更新时间从 10 改成 6:
B: 9, C: 8, D: 7, A: 6
如果查询视图已经刷新,后续结果集的排序边界发生了变化。根据具体变化和游标位置,可能出现重复或跳过。
情况三:删除尚未访问的文档
如果 C 在第二次请求前被删除,后续查询自然不会再返回它。若业务把分页结果理解为初始时刻的完整快照,这也是一种不一致。
这里的关键不是 search_after 算错了,而是它每次都在当前可见索引视图上继续执行。它没有保存“第一页请求时的文档集合”。
五、PIT:固定一个 Point in Time 查询视图
1. PIT 的定义
PIT(Point in Time)可以理解为:
为一个或多个索引创建一个可持续使用的搜索视图,使后续搜索请求基于创建 PIT 时的索引状态读取。
PIT 不是数据库事务,也不提供写入锁。创建 PIT 后,其他请求仍然可以写入、更新和删除索引。只是这些变化不会出现在已经创建的 PIT 搜索视图中。
PIT 主要解决的是:
search_after:我从哪里继续?
PIT:我在什么数据视图上继续?
两者通常配合使用。
2. PIT 的生命周期
一个典型生命周期如下:
- 调用 Open PIT;
- Elasticsearch 为相关索引创建搜索上下文;
- 使用 PIT ID 发起多次
_search; - 每次请求可以续期
keep_alive; - 导出结束、失败或取消时关闭 PIT;
- 如果超过保留时间,PIT 可能过期,后续请求失败。
创建 PIT:
curl -X POST 'http://localhost:9200/products/_pit?keep_alive=2m'
可能返回:
{
"id": "46ToAw..."
}
使用 PIT 时,查询请求通过 pit 指定 PIT,而不是再次指定索引:
curl -X POST 'http://localhost:9200/_search' \
-H 'Content-Type: application/json' \
-d '{
"size": 2,
"pit": {
"id": "46ToAw...",
"keep_alive": "2m"
},
"query": {
"match_all": {}
},
"sort": [
{
"updated_at": "desc"
},
{
"product_id": "asc"
}
]
}'
后续请求必须使用响应中最新返回的 PIT ID。PIT ID 是不透明值,不应由程序自行解析或拼接。
关闭 PIT:
curl -X DELETE 'http://localhost:9200/_pit' \
-H 'Content-Type: application/json' \
-d '{
"id": "46ToAw..."
}'
关闭操作应放在 finally 或等价的清理逻辑中,不能只依赖自然过期。
3. PIT 与 search_after 的完整示例
下面假设索引映射中有:
updated_at:日期类型;product_id:唯一的keyword字段;name:可搜索文本。
第一页请求:
curl -X POST 'http://localhost:9200/_search' \
-H 'Content-Type: application/json' \
-d '{
"size": 2,
"track_total_hits": false,
"pit": {
"id": "46ToAw...",
"keep_alive": "2m"
},
"query": {
"match": {
"name": "keyboard"
}
},
"sort": [
{
"updated_at": "desc"
},
{
"product_id": "asc"
}
]
}'
假设最后一条命中返回:
{
"_id": "p101",
"_source": {
"product_id": "p101",
"name": "mechanical keyboard"
},
"sort": [
"2025-01-10T11:00:00.000Z",
"p101"
]
}
第二页:
curl -X POST 'http://localhost:9200/_search' \
-H 'Content-Type: application/json' \
-d '{
"size": 2,
"track_total_hits": false,
"pit": {
"id": "46ToAw...",
"keep_alive": "2m"
},
"query": {
"match": {
"name": "keyboard"
}
},
"sort": [
{
"updated_at": "desc"
},
{
"product_id": "asc"
}
],
"search_after": [
"2025-01-10T11:00:00.000Z",
"p101"
]
}'
读取逻辑是:
打开 PIT
last_sort = null
循环:
使用 PIT 查询一页
如果 hits 为空:
结束
处理 hits
last_sort = 最后一条 hit.sort
下一次请求携带 search_after = last_sort
无论成功、失败还是取消:
关闭 PIT
track_total_hits: false 表示不要求每一页都精确统计匹配总数。在大量结果的导出场景中,这通常可以减少无关成本;如果界面必须显示精确总数,则应明确承担统计成本,或将统计请求与翻页请求分离。
4. PIT 的隐式排序辅助字段
使用 PIT 搜索时,Elasticsearch 会使用一个内部排序辅助值作为并列文档的 tie-breaker。响应中的 sort 数组可能因此比请求中显式声明的排序字段多一个值。
实际代码必须:
- 完整保存最后一条命中的
sort数组; - 原样传回
search_after; - 不要只截取自己认识的前几个值;
- 不要手工构造或修改内部 tie-breaker。
官方语义中,PIT 场景常使用 _shard_doc 作为内部文档顺序辅助字段。它适合配合 PIT 使用,但不应被当作没有 PIT 时的全局稳定业务排序键。
如果只需要遍历 PIT 中的全部文档、不需要业务排序,可以考虑 _shard_doc 排序,并在不需要精确总数时关闭 track_total_hits,这是 Elasticsearch 针对 PIT 全量扫描提供的高效路径之一。但这仍然不是业务主键排序,也不提供跨 PIT 的断点语义。
六、PIT 提供了什么一致性,又没有提供什么
1. PIT 固定的是搜索视图
创建 PIT 后:
t0:创建 PIT
t1:写入文档 X
t2:更新文档 A
t3:删除文档 B
t4:使用 PIT 查询
在正常语义下,t4 的 PIT 查询仍基于 t0 建立的视图:
X不会因为后来写入而出现在这个 PIT 中;A返回 PIT 视图中可见的旧版本;B如果在创建 PIT 时可见,则仍可能被返回;- 后续 refresh 不会把新版本自动加入这个 PIT。
这让多次分页请求可以基于同一个逻辑视图执行。
2. PIT 不是跨请求写入事务
PIT 不会:
- 阻止其他客户端写入;
- 阻止文档更新;
- 阻止文档删除;
- 保证导出文件一次性提交;
- 保证多个独立 PIT 之间看到同一个时间点;
- 替代数据库事务或快照备份。
如果要求“导出开始前的精确全集”,PIT 通常可以提供查询层面的时间点视图。但如果导出跨越多个独立索引、多个 PIT,或者需要与外部数据库事务严格对齐,就必须额外定义边界。
3. PIT 的资源代价
PIT 会保留搜索所需的索引读取上下文,使相关 Lucene 资源不能立即按普通查询后的路径释放。长时间存活、数量过多的 PIT 可能导致:
- 文件句柄和磁盘资源保持时间变长;
- segment merge 后旧 segment 不能及时清理;
- 节点搜索上下文数量增加;
- 节点负载和故障恢复压力上升。
keep_alive 不是“查询必须在多久内完成”的超时,而是搜索上下文允许保持多久。每次成功请求可以续期,但它不应被设置成远大于实际处理时间的值。
七、Scroll:保存服务端搜索上下文的批量读取方式
1. Scroll 的工作方式
Scroll 的核心流程是:
- 第一次
_search创建 scroll 搜索上下文; - 服务端保存这个上下文及当前读取位置;
- 客户端拿到
_scroll_id; - 后续请求只提交 scroll ID,继续取得下一批;
- 直到返回空结果;
- 显式清理 scroll 上下文。
创建 scroll:
curl -X POST 'http://localhost:9200/products/_search?scroll=1m' \
-H 'Content-Type: application/json' \
-d '{
"size": 1000,
"query": {
"match_all": {}
},
"sort": [
"_doc"
]
}'
响应中会有:
{
"_scroll_id": "DnF1ZXJ5VGhlbkZldGNo...",
"hits": {
"hits": [
"...最多 1000 条..."
]
}
}
继续读取:
curl -X POST 'http://localhost:9200/_search/scroll' \
-H 'Content-Type: application/json' \
-d '{
"scroll": "1m",
"scroll_id": "DnF1ZXJ5VGhlbkZldGNo..."
}'
每次响应都可能返回新的 _scroll_id,客户端必须使用最新的值,而不能永久复用第一次返回的 ID。
清理:
curl -X DELETE 'http://localhost:9200/_search/scroll' \
-H 'Content-Type: application/json' \
-d '{
"scroll_id": [
"DnF1ZXJ5VGhlbkZldGNo..."
]
}'
也可以清理多个 scroll ID。程序应在正常结束、异常退出可控处理的路径中主动清理。
2. Scroll 的一致性边界
Scroll 搜索上下文会保存初始搜索时的读取视图。后续 scroll 请求继续读取这一上下文,而不是每次重新从当前索引视图开始。
因此,Scroll 适合把一次批量读取理解为:
在开始时建立一个搜索视图,
之后沿着这个视图分批消费结果。
但要注意以下边界:
- 它不是数据库事务;
- 它不锁住索引;
- 它不阻止写入;
- 它不保证下游文件写入成功;
- scroll 上下文过期后,不能从原上下文继续;
- 分片不可用、节点故障或上下文被清理时,读取可能失败。
Scroll 的存活时间是服务端上下文生命周期,不是客户端可以无限暂停的许可。每次请求可以携带新的 scroll 时间,但如果客户端处理一批数据的时间超过上下文保留时间,下一次请求仍可能失败。
3. Scroll 与排序
对于大规模批量遍历,如果不要求业务排序,通常使用:
"sort": ["_doc"]
这样可以避免为全局业务排序保留额外候选,适合“尽快把所有匹配文档读出来”。
如果需要按业务字段排序,Scroll 仍然可以工作,但排序成本可能显著增加。此时应确认:
- 排序字段有合适的 doc values;
- 排序是否真的必要;
- 是否可以使用 PIT +
search_after; - 是否需要使用 sliced scroll 并行处理。
Scroll 不使用 from 逐页跳转,也不使用 search_after 作为主要游标。它的游标是服务端返回的 _scroll_id。
八、search_after、PIT 和 Scroll 的关系
可以用下面的表格区分它们:
| 机制 | 保存什么状态 | 是否固定查询视图 | 是否适合随机跳页 | 典型用途 |
|---|---|---|---|---|
from + size |
基本不保存游标 | 否 | 是,但深页代价高 | 浅分页、页码跳转 |
search_after |
客户端保存排序游标 | 否 | 否 | 连续翻页、实时列表 |
PIT + search_after |
PIT 保存视图,客户端保存排序游标 | 是 | 否 | 稳定分页、长结果集 |
| Scroll | 服务端保存搜索上下文和读取位置 | 是,基于初始搜索视图 | 否 | 批量读取、传统导出 |
最常见的选择是:
- 用户列表:
search_after,是否加 PIT 取决于是否需要跨页稳定视图; - 稳定的长列表:PIT +
search_after; - 一次性批量导出:Scroll 或 PIT +
search_after; - 需要快速遍历、不要求业务排序:Scroll +
_doc,或 PIT +_shard_doc; - 需要页码跳转:浅分页使用
from + size,不要试图用search_after模拟任意页码。
九、导出不是“循环调用搜索”这么简单
导出至少要定义四个问题:
- 导出开始时看到哪个数据视图?
- 每批的游标是什么?
- 失败后从哪里恢复?
- 输出中重复一条记录是否可接受?
1. 导出方案一:Scroll
Scroll 适合传统批量导出:
创建 scroll
读取一批
写入输出
读取下一批
直到为空
清理 scroll
典型适用场景:
- 导出一个相对固定的索引内容;
- 结果不要求按业务字段排序;
- 导出程序可以在一次上下文生命周期内完成;
- 下游更关注吞吐而不是用户可见的分页体验。
批次大小需要根据 _source 大小、网络带宽、下游写入速度和节点内存调节。批次不是越大越好:
- 太小:请求次数多,协议和调度开销高;
- 太大:单次响应大,GC、网络重试和下游缓冲压力增加。
2. 导出方案二:PIT + search_after
对于希望显式控制排序、游标和生命周期的导出,可以使用 PIT + search_after:
打开 PIT
使用固定 query 和 sort 查询第一页
保存最后一条 hit.sort
处理当前页
下一页携带 search_after
每次续期 PIT
导出结束后关闭 PIT
优点:
- 排序逻辑清楚;
- 游标由客户端掌握;
- 查询请求本身更接近普通
_search; - 可以按照业务排序稳定遍历;
- 适合把游标写入外部进度记录。
但它不等于自动支持断点续传。PIT 过期后,原来的 search_after 只是一组排序值,不能保证在新 PIT 中仍代表同一个数据视图。
十、导出失败时,为什么会重复或丢失
设导出每批 1000 条,处理第 5 批时,程序要做两件事:
- 把数据写入文件或消息系统;
- 保存“下一次从哪里开始”的游标。
这两步通常不能天然组成一个跨系统原子事务。
顺序一:先写数据,再保存游标
写入第 5 批
程序崩溃
游标仍停在第 4 批
重启后重新读取第 5 批
结果:可能重复。
顺序二:先保存游标,再写数据
保存第 5 批之后的游标
程序随后崩溃
第 5 批尚未写入
重启后从第 6 批继续
结果:可能丢失。
因此,导出系统必须选择一种额外策略:
方案 A:下游幂等写入
给每条记录设置稳定业务键,例如:
导出任务 ID + 文档唯一业务 ID
下游以该键去重或覆盖,允许重试造成重复发送,但最终结果不重复。
方案 B:按批次写入临时文件,再原子提交
一个批次完整写入临时文件并 fsync 或完成对象存储上传后,才记录该批次已提交。恢复时只认已提交批次。
方案 C:接受至少一次语义
明确规定导出是 at-least-once,允许重复,由下游清洗。这个方案实现简单,但必须在接口和数据契约中说明。
方案 D:固定外部数据版本
如果 Elasticsearch 数据来自带版本或时间边界的数据管道,可以把导出范围限定为:
ingest_version <= V
这样即使重新创建 PIT,也能通过业务条件重新确定同一个逻辑集合。但这个能力来自业务数据模型,不是 PIT 自动提供的。
十一、PIT 过期和 Scroll 失效如何处理
1. PIT 过期
PIT 过期、节点异常或相关上下文不可用时,搜索请求可能返回上下文不存在或搜索失败类错误。
不能简单地做:
PIT 失败
重新打开 PIT
继续使用旧 search_after
因为新 PIT 对应的是新的索引视图。旧游标可能导致:
- 新数据被插入游标之前;
- 文档更新后排序值改变;
- 旧视图中存在的文档在新视图中已删除;
- 结果集合与旧 PIT 不再相同。
正确处理取决于业务要求:
- 如果只是实时列表:重新开始或回退到一个明确的业务边界;
- 如果是允许重复的导出:重新打开 PIT,并由下游按业务键去重;
- 如果要求严格无漏导出:需要更长但合理的 PIT 生命周期、降低单批处理时间,或者使用外部版本边界和可重放机制;
- 如果要求强一致快照:应重新评估 Elasticsearch 查询快照是否足够,不能把 PIT 当作备份快照或数据库事务。
2. Scroll 过期
Scroll ID 失效后,通常不能从服务端上下文中恢复原来的位置。若程序只保存了 scroll ID,而没有保存已经成功输出的业务键或批次信息,恢复能力会比较弱。
因此,批处理导出即使采用 Scroll,也应记录:
- 导出任务标识;
- 已提交批次;
- 每条数据的幂等键;
- 错误原因和重试次数;
- 最后成功处理时间。
Scroll 解决读取上下文保存,不解决导出事务。
十二、并发导出与 sliced scroll
1. 为什么要切片
单个 scroll 或单个 PIT 游标是串行消费路径。如果数据量很大,可以把查询拆成多个 slice,让多个工作线程或任务并行读取。
概念上:
slice 0:读取属于第 0 片的数据
slice 1:读取属于第 1 片的数据
slice 2:读取属于第 2 片的数据
...
Scroll 的 sliced scroll 通常在查询中加入:
"slice": {
"id": 0,
"max": 4
}
每个 slice 使用独立的 scroll 上下文和 scroll ID。其他 slice 依次使用不同的 id,但共享相同的 max。
2. 并行切片的代价
切片不是免费并行化:
- 每个 slice 都会创建搜索上下文;
- 每个 slice 都会占用请求、网络和下游处理资源;
- slice 数过多可能压垮分片和协调节点;
- 某些查询并不适合简单切片;
- 每个 slice 的输出顺序通常不能直接合并为全局业务顺序。
因此,切片更适合:
无须全局排序的批量导出
如果要求导出文件严格按照 updated_at 全局排序,多个 slice 各自并行输出后,必须额外做归并排序,这会重新引入缓冲和复杂性。
十三、一个可运行的 PIT + search_after 导出骨架
下面用 Python 标准库调用 Elasticsearch HTTP API,展示核心生命周期。示例假设:
- Elasticsearch 地址为
http://localhost:9200; - 索引为
products; updated_at可排序;product_id是全局唯一keyword字段;- 导出只读取
_source; - 输出文件允许通过临时文件完成最终替换。
import json
import os
import tempfile
import urllib.request
import urllib.error
ES = "http://localhost:9200"
INDEX = "products"
KEEP_ALIVE = "2m"
PAGE_SIZE = 500
def request(method, path, body=None):
data = None
headers = {}
if body is not None:
data = json.dumps(body).encode("utf-8")
headers["Content-Type"] = "application/json"
req = urllib.request.Request(
ES + path,
data=data,
headers=headers,
method=method,
)
with urllib.request.urlopen(req, timeout=30) as resp:
return json.loads(resp.read().decode("utf-8"))
def export_products(target_path):
pit = None
temp_path = None
try:
# 1. 打开 PIT。响应的 id 是不透明值。
opened = request(
"POST",
f"/{INDEX}/_pit?keep_alive={KEEP_ALIVE}",
)
pit = opened["id"]
fd, temp_path = tempfile.mkstemp(
prefix="products-",
suffix=".jsonl",
)
with os.fdopen(fd, "w", encoding="utf-8") as output:
search_after = None
while True:
body = {
"size": PAGE_SIZE,
"track_total_hits": False,
"pit": {
"id": pit,
"keep_alive": KEEP_ALIVE,
},
"query": {
"match_all": {}
},
"sort": [
{"updated_at": "desc"},
{"product_id": "asc"},
],
}
if search_after is not None:
body["search_after"] = search_after
page = request("POST", "/_search", body)
# PIT ID 可能在响应中更新,必须使用最新值。
if "pit_id" in page:
pit = page["pit_id"]
hits = page["hits"]["hits"]
if not hits:
break
for hit in hits:
record = {
"_id": hit["_id"],
"_source": hit.get("_source", {}),
}
output.write(json.dumps(
record,
ensure_ascii=False,
) + "\n")
# 必须保存最后一条命中的完整 sort 数组。
search_after = hits[-1]["sort"]
# 2. 临时文件完整写完后,再替换最终文件。
os.replace(temp_path, target_path)
temp_path = None
except urllib.error.HTTPError as exc:
# 生产代码应记录响应体、任务 ID、PIT 状态和最后成功批次。
raise RuntimeError(
f"Elasticsearch request failed: {exc.code}"
) from exc
finally:
if pit is not None:
try:
request("DELETE", "/_pit", {"id": pit})
except Exception:
# 清理失败应告警,但不能覆盖原始导出异常。
pass
if temp_path is not None:
try:
os.remove(temp_path)
except FileNotFoundError:
pass
这个示例的关键点不在 Python 语法,而在状态顺序:
- PIT ID 是服务端返回的状态;
search_after来自上一页最后一条命中的完整sort;- 每页请求都续期 PIT;
- 结果先写临时文件;
- 全部完成后才替换最终文件;
- 任意异常都尝试关闭 PIT;
- PIT 失效时直接报错,而不是假装新建 PIT 后可以无条件续接旧游标。
这个示例仍然没有把“写文件”和“记录导出任务状态”变成跨系统事务。如果输出是数据库、对象存储或消息队列,需要按照对应系统设计幂等和提交协议。
十四、导出边界:文档更新、删除和 _source
PIT 或 Scroll 固定的是搜索读取视图,但 Elasticsearch 的更新通常是新版本文档写入和旧版本删除的组合。从查询者角度看,PIT 会读取创建时可见的版本。
因此,在 PIT 中导出文档时:
- 文档创建后才写入:不会出现在已创建的 PIT 中;
- 文档在 PIT 创建后更新:通常仍读取旧版本;
- 文档在 PIT 创建后删除:旧视图中仍可能可见;
_source是该视图中命中的文档源内容,不是导出时刻最新内容。
这对“导出当前最新状态”很重要:PIT 导出的可能是一个一致的历史视图,而不是导出结束时的最新数据。
如果业务要求的是“截至某个业务时间点的数据”,应在查询条件中显式使用业务时间或版本字段,而不能只依赖 Elasticsearch 的 PIT 创建时间。尤其是当写入链路存在延迟、refresh 延迟或跨系统事件顺序时,索引可见时间和业务发生时间并不相同。
十五、排序字段的实际要求
稳定深分页的排序设计通常要满足以下条件。
1. 排序字段可排序
文本字段不能直接按分析后的 token 进行普通业务排序。通常需要使用 keyword 子字段:
"sort": [
{
"customer_name.keyword": "asc"
}
]
日期、数值和 keyword 字段通常更适合作为排序字段。
2. 必须处理并列值
仅按时间排序:
"sort": [
{
"created_at": "asc"
}
]
如果同一毫秒创建了很多文档,排序边界不完整。应追加唯一键:
"sort": [
{
"created_at": "asc"
},
{
"event_id": "asc"
}
]
3. 字段值类型和格式要稳定
游标数组中的值必须和排序字段的值语义一致。不要在第一页使用一种日期格式、第二页又改变字段类型或排序方向。
4. 多索引查询需要统一映射
如果查询 logs-2024-* 和 logs-2025-*,所有参与排序的字段都应具有兼容映射。否则可能出现字段类型冲突、排序失败或结果语义不一致。
5. 缺失值必须有明确策略
某些文档没有排序字段时,Elasticsearch 会按照排序配置处理缺失值。导出程序应确认缺失值在 sort 数组中的表现,并避免把不同类型的游标值混入外部状态。
十六、失败表现和诊断方法
1. Result window is too large
常见表现是请求超过默认的 from + size 窗口限制。
诊断重点:
- 查看请求中的
from和size; - 确认是否误把深分页实现成了页码分页;
- 不要先急着调大
index.max_result_window; - 改用
search_after或 Scroll。
2. search_after 返回重复或漏数据
优先检查:
- 是否没有使用 PIT;
- 分页期间是否发生 refresh、写入、更新或删除;
- 是否缺少唯一 tie-breaker;
- 是否错误地使用了当前页第一条而不是最后一条的
sort; - 下一页是否修改了 query、sort、过滤条件或索引范围;
- 是否截断了 PIT 响应中的完整
sort数组; - 是否把不同 PIT 的游标混用了。
3. PIT 或 Scroll 上下文不存在
常见原因:
keep_alive太短;- 两次请求间处理时间过长;
- 节点重启或故障;
- 搜索上下文被清理;
- 客户端没有使用响应返回的最新 PIT ID 或 scroll ID;
- 请求发送到了不再持有相关状态的环境,或集群发生了影响上下文的故障。
处理时应区分:
- 可以从头重试的实时查询;
- 允许重复、可幂等恢复的导出;
- 要求严格无漏的导出。
三者的恢复策略不同。
4. 节点资源异常
应观察:
- 搜索线程池拒绝;
- JVM 堆和 GC;
- 文件句柄;
- segment 数量和 merge 状态;
- 活跃 PIT、scroll 上下文数量;
- 查询响应时间和网络响应大小;
- 导出程序自身的下游写入积压。
长时间保留 PIT 或同时开启大量 scroll,可能让“查询没报错”变成“集群资源逐步恶化”。这类问题通常需要从生命周期、并发度和批次大小一起诊断。
十七、生产取舍
用户列表
如果用户只需要连续点击“下一页”:
search_after
通常比深层 from + size 更合适。
如果用户要求在一段浏览期间结果不随刷新变化:
PIT + search_after
更合适,但应设置合理的 PIT 过期时间,并在会话结束时关闭。
如果用户需要直接跳到第 50 页:
浅页使用 from + size
深页不要强行模拟页码
可以改成按条件筛选、时间范围、游标分页或搜索条件缩小结果集。
批量导出
如果不要求业务排序、追求吞吐:
Scroll + _doc
是典型方案。
如果要求明确的业务排序、客户端掌握游标,或者需要更清晰地控制分页请求:
PIT + search_after
更容易表达查询语义。
无论选择哪一种,都需要单独设计:
- 导出任务状态;
- 批次提交;
- 重试策略;
- 下游幂等键;
- 上下文清理;
- 上下文失效后的恢复策略。
数据量极大或要求长期断点续传
PIT 和 Scroll 都不是永久快照。若导出可能持续很长时间,或必须在数小时、数天后从精确位置恢复,应考虑:
- 按业务时间或版本分段导出;
- 使用可重放的索引快照或离线副本;
- 将唯一版本号纳入查询条件;
- 让下游以业务主键幂等写入;
- 将任务拆成可独立提交的区间。
这里的核心原则是:查询游标只能描述“怎么继续查”,不能自动承担“跨故障保存完整数据集”的责任。
十八、最后的判断框架
可以按三个问题选择方案。
问题一:是否需要深度访问?
- 不需要:
from + size; - 需要连续向后读取:
search_after或 Scroll; - 需要任意页码:重新设计交互或查询边界。
问题二:是否需要跨请求稳定视图?
- 不需要,允许读取实时变化:
search_after; - 需要同一时间点视图:PIT +
search_after,或 Scroll。
问题三:是否是导出而不是交互分页?
- 批量、无业务排序:Scroll;
- 有业务排序、希望客户端管理游标:PIT +
search_after; - 要求故障后严格无漏恢复:在上述方案之外增加版本边界、幂等输出或可重放机制。
search_after 解决游标推进,PIT 解决跨请求视图稳定,Scroll 解决服务端批量读取上下文,导出则必须继续解决输出提交和故障恢复。只有把这四个层次分开,才能正确判断一次深分页或导出设计到底保证了什么。
系列导航与关联阅读
- 系列入口:数据库完整学习路线:从关系模型、事务索引到分布式与向量检索
- 上一篇:Elasticsearch 相关性调优:BM25、Boost、Function Score 和评测集
- 下一篇:Elasticsearch 安全与多租户:TLS、角色、文档权限、审计和隔离
- 延伸:Elasticsearch 查询与聚合:Query DSL、相关性、分页和统计
- 延伸:Elasticsearch 数据写入与运维:Bulk、Ingest、ILM、快照和升级
官方资料
本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。

评论
0 条讨论