数据库基础体系 · 第 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 个主分片:

  1. 协调节点把查询发送到 3 个分片;
  2. 每个分片都需要找出本地排序靠前的 from + size = 1020 条候选;
  3. 协调节点合并各分片候选;
  4. 再跳过前 1000 条,返回 20 条。

因此,页码越深,每个分片需要保留和排序的候选越多。粗略地说,协调过程的候选规模与下面的量相关:

O(S×(F+Z))O(S \times (F + Z))

其中:

  • SS:分片数量;
  • FFfrom
  • ZZsize

这不是精确的内存复杂度公式,因为实际执行还取决于查询类型、排序字段、缓存和 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 返回的每一页都没有遗漏,写入文件或下游系统时仍然可能失败:

  1. 第 10 页已经写入文件;
  2. 程序在保存进度前崩溃;
  3. 重启后从第 10 页重新读取;
  4. 文件中出现重复数据。

这已经不是 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 条,因为它无法预先知道其他分片会提供哪些更靠前的文档。

这解释了两件事:

  1. 深分页的成本主要发生在“找出前面那些不返回给用户的结果”;
  2. 仅仅把 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 的排序条件

假设排序键是:

K(d)=(t(d),id(d))K(d) = (t(d), id(d))

其中:

  • t(d)t(d):文档的更新时间,按降序;
  • id(d)id(d):业务唯一标识,按升序。

对于当前游标 KcK_c,下一页应返回满足下列条件的文档:

t(d)<tct(d) < t_c

或者:

t(d)=tcid(d)>idct(d) = t_c \land id(d) > id_c

这就是多字段排序的字典序。第二个字段的作用是打破第一个字段的并列。

因此,稳定分页通常需要一个真正唯一的排序键。例如:

"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 的生命周期

一个典型生命周期如下:

  1. 调用 Open PIT;
  2. Elasticsearch 为相关索引创建搜索上下文;
  3. 使用 PIT ID 发起多次 _search
  4. 每次请求可以续期 keep_alive
  5. 导出结束、失败或取消时关闭 PIT;
  6. 如果超过保留时间,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 的核心流程是:

  1. 第一次 _search 创建 scroll 搜索上下文;
  2. 服务端保存这个上下文及当前读取位置;
  3. 客户端拿到 _scroll_id
  4. 后续请求只提交 scroll ID,继续取得下一批;
  5. 直到返回空结果;
  6. 显式清理 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. 导出开始时看到哪个数据视图?
  2. 每批的游标是什么?
  3. 失败后从哪里恢复?
  4. 输出中重复一条记录是否可接受?

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 批时,程序要做两件事:

  1. 把数据写入文件或消息系统;
  2. 保存“下一次从哪里开始”的游标。

这两步通常不能天然组成一个跨系统原子事务。

顺序一:先写数据,再保存游标

写入第 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 语法,而在状态顺序:

  1. PIT ID 是服务端返回的状态;
  2. search_after 来自上一页最后一条命中的完整 sort
  3. 每页请求都续期 PIT;
  4. 结果先写临时文件;
  5. 全部完成后才替换最终文件;
  6. 任意异常都尝试关闭 PIT;
  7. 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 窗口限制。

诊断重点:

  • 查看请求中的 fromsize
  • 确认是否误把深分页实现成了页码分页;
  • 不要先急着调大 index.max_result_window
  • 改用 search_after 或 Scroll。

2. search_after 返回重复或漏数据

优先检查:

  1. 是否没有使用 PIT;
  2. 分页期间是否发生 refresh、写入、更新或删除;
  3. 是否缺少唯一 tie-breaker;
  4. 是否错误地使用了当前页第一条而不是最后一条的 sort
  5. 下一页是否修改了 query、sort、过滤条件或索引范围;
  6. 是否截断了 PIT 响应中的完整 sort 数组;
  7. 是否把不同 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 解决服务端批量读取上下文,导出则必须继续解决输出提交和故障恢复。只有把这四个层次分开,才能正确判断一次深分页或导出设计到底保证了什么。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。