当前位置: 首页 > news >正文

Scroll 翻页查询


一、 原理

Scroll 的工作方式和 from + size 完全不同。它不排序

  1. 首次请求生成一个Lucene Scroll ID,同时在协调节点上持有一个ScrollContext,保存了本次查询涉及的所有分片的游标状态
  2. 后续每个 scroll 请求都携带这个 ID,协调节点到每个分片上继续往后取数据,每个分片记住自己扫描到了哪里
  3. 每次返回size条给客户端,直到所有分片数据取完

关键区别:Scroll 使用的是非活跃快照。发起 Scroll 请求的那一刻,Lucene 在段级别做了一个引用计数快照。后续即使有新写入、段合并甚至删除操作,这个 Scroll 看到的始终是快照时刻的数据。代价是快照期间消耗额外的磁盘空间,直到 Scroll 超时释放。

每个分片内的遍历是直接走 doc values 的顺序扫描_doc或指定排序字段),不像 from + size 那样做全局 Top-K 归并。所以内存开销是常数级,性能稳定。


二、Scroll 占用的额外磁盘空间是什么

Scroll 不拷贝文档内容。它做的是阻止 Lucene 段合并清理旧段,通过段引用计数来实现。

Lucene 的段是不可变的。删除和更新操作不会原地修改段,而是在新段中标记删除,旧段在合并时被物理清理。Scroll 创建快照时,给当前所有相关段的引用计数加 1。只要 Scroll 没超时释放,这些段——以及其中被标记为删除的文档——就不能被合并线程回收。

所以额外消耗的磁盘空间,等于 Scroll 存活期间本应被段合并回收的废弃数据量。

示例计算:

原始 10GB Scroll 期间:删 2GB + 更新 1GB 文档 + 写入 1GB 新数据 磁盘占用 ≈ 8 + 2(删除锁住)+ 1(旧版本锁住)+ 1(新版本)+ 1(新数据)= 13GB 额外消耗 = 2 + 1 = 3GB

Scroll 释放后段合并回收 3GB,磁盘回到约 10GB。


三、 更新、删除、新增操作在 Scroll 期间的行为

更新和写入的数据走的是完全不同的路径,和 Scroll 持有的旧段互不干扰。

写入新数据

新文档写入 Lucene 时会分配全新的段,和 Scroll 引用的旧段没有任何关系。

磁盘上: 旧段(Segment 1-5) ← Scroll 持有引用,冻结 新段(Segment 6-7) ← 新写入的数据落在这里,Scroll 看不到

Scroll 的快照不包含新段,所以新写入的文档在 Scroll 遍历期间完全不可见。磁盘上就是正常的增量——新段占多少就是多少,Scroll 不影响新数据的存储开销。

更新已有数据

更新在 Lucene 中的实现是「标记删除旧文档 + 写入新文档」,不是原地修改。

更新前: Segment 1: [doc_A v1] ← Scroll 持有引用 更新 doc_A 后: Segment 1: [doc_A v1](标记删除) ← Scroll 还抓着,无法合并回收 Segment 8: [doc_A v2] ← 新版本,Scroll 看不到

关键影响在于:v1 和 v2 同时在磁盘上

  • 正常情况下,Segment 1 会在合并时回收 doc_A v1 的空间,文档只存一份 v2
  • Scroll 存活期间,v1 所在的段被引用锁住无法合并,v1 和 v2 同时占用磁盘

所以更新操作在 Scroll 期间会造成临时性的双倍存储——旧版本被 Scroll 锁住不释放,新版本照常落盘。这个额外消耗的上限是 Scroll 存活期内被更新文档的旧版本总大小。

总结三种操作的磁盘影响

操作对 Scroll 可见性磁盘额外消耗
删除被删文档仍可见被删文档空间无法释放
更新看到旧版本旧版本 + 新版本同时存在
新增不可见仅正常的新段存储,无额外消耗

四、 Scroll 存活期间普通查询会不会返回新旧两种不同版本数据

不会。普通查询不会因为 Scroll 还锁着旧段就意外返回旧版本数据。

关键机制:live docs bitset

Lucene 中每个段维护一个位图(bitset),每一位对应段内的一个文档:1 表示存活,0 表示已删除。删除操作本质上是把这个位翻转,而不是物理删除数据。

更新也一样——先在旧段中把旧文档标记为删除(位翻转为 0),然后在新段中写入新版本。

更新前: Segment 1 live docs: [1, 1, 1, 1] ← doc_A 存活 更新 doc_A 后: Segment 1 live docs: [0, 1, 1, 1] ← doc_A 标记删除,位翻为 0 Segment 8 新文档 doc_A v2 ← 新版本写在这里

普通查询和 Scroll 查询的关键差异

普通查询使用的是当前最新的 IndexReader。打开新的 IndexReader 时,它会读取每个段当前的live docs bitset。Segment 1 中 doc_A 的位已经是 0,查询时直接被过滤掉,不会被返回。即使 Segment 1 因为 Scroll 还锁着不能被合并清理,位图已经告诉查询引擎「这条文档是已删除的」。

Scroll 查询使用的是创建时的 IndexReader。它从打开那一刻起就持有自己的 live docs 快照视图。当时的 live docs 显示 doc_A 位为 1(还没被删除/更新),所以 Scroll 遍历时会看到旧版本。

用图来对比

磁盘上同时存在: Segment 1: [doc_A v1] live doc bit = 0(已删除) Segment 8: [doc_A v2] live doc bit = 1(存活) 普通查询 → IndexReader 看到 live docs 当前状态 → Segment 1 中 doc_A 被过滤掉 → 返回 Segment 8 中的 v2 ✓ Scroll → IndexReader 持有自己的快照 live docs → Segment 1 中 doc_A 位仍为 1 → 返回 Segment 1 中的 v1

所以实际上同一个段、同一份数据,同一个物理文件,两个查询看到的结果不同——不是因为数据有两份,而是因为 live docs bitset 的视图不同。

唯一的例外:refresh 窗口

在 refresh 间隔内(默认 1 秒),新段还没被打开,普通查询的 IndexReader 还不知道 Segment 8 的存在。此时旧版本已被标记删除,新版本尚未可见,查询 doc_A 会返回空——但这和 Scroll 锁不锁旧段没有关系,即使没有 Scroll 存在,这个窗口内也一样查不到。


五、NRT(近实时)搜索模型与查询一致性

refresh 机制

写入 → 内存 buffer → refresh → 生成新 segment → 搜索可见,默认间隔 1 秒。

t = 0.0s: doc = "旧值" t = 0.1s: 更新为 "新值",写入内存 buffer t = 0.2s: 查询 → "旧值"(还没 refresh) t = 1.1s: refresh 触发 t = 1.2s: 查询 → "新值"

Read Your Own Writes 的三种方案

  1. 实时 GETGET /index/_doc/id直接从 translog 检索,绕过 refresh
  2. refresh=wait_forPUT /index/_doc/1?refresh=wait_for阻塞写入直到 refresh 完成
  3. preference=_primary:强制走主分片,避免读到未同步的副本

Scroll 的"不一致"是刻意冻结快照时刻视图——整个存活期内(数分钟甚至更长)都看不到任何新数据。普通查询的"不一致"是 NRT 延迟窗口(默认 1 秒),写入最终都会可见。一个是设计特性,一个是时间窗口问题。


六、Scroll 与普通查询的差异总览

维度Scroll普通查询
视图创建时刻的快照,冻结不变每次查询看到最新 refresh 后的数据
对新数据可见性不可见refresh 后可见
对旧版本可见性如果快照时未被删除则可见按当前 live docs bitset 过滤
内存模型常数级from+size 线性增长
排序不排序,顺序扫描全局 Top-K 排序
跳页不支持支持(但有深度限制)
聚合不支持支持
磁盘影响阻止废弃段回收无影响

七、 Scroll 存在意义和使用场景

Scroll 使用非活跃快照:发起时刻 Lucene 在段级别做引用计数快照,后续写入、删除、段合并都不影响 Scroll 看到的视图。

代价:快照期间阻止段合并回收空间。

ES 7.x 以后官方不建议用 Scroll 做实时分页,建议迁移到 search_after。Scroll 现在的定位是一次性全量遍历

适用场景

  • 全量数据导出(百万、千万、亿级),生成报表、同步到数据仓库
  • 定时批处理任务,对实时性无要求
  • 重建索引(reindex)过程中的源端读取

注意事项

  • keep_alive 尽可能短——不是为了释放 Scroll 对象本身,而是为了尽快让段合并回收空间
  • Scroll 不支持聚合、不支持跳页
  • 单次 size 受 max_result_window 限制,但游标可无限翻

八、测试脚本

#!/usr/bin/env python3""" Easysearch Scroll 分页测试脚本 覆盖内容:1. 创建测试索引并写入大量数据2. Scroll 全量遍历3. Scroll 存活期间写入/更新/删除行为验证 使用方式: 修改下方 EASYSEARCH_HOST、USERNAME、PASSWORD 后执行 python3 easysearch_scroll_test.py"""importtimeimportjsonimportrandomimportstring# ============================================================# 配置区# ============================================================EASYSEARCH_HOST="https://localhost:9200"# 修改为你的 Easysearch 地址USERNAME="admin"# 如果不需要认证,置空PASSWORD=""# 如果不需要认证,置空VERIFY_CERTS=False# 自签名证书时设为 FalseINDEX_NAME="scroll_test_demo"DOC_COUNT=20000# 写入文档总数SCROLL_SIZE=1000# 每批取多少条SCROLL_KEEP="2m"# Scroll 存活时间# ============================================================# 初始化客户端# ============================================================def create_client():"""创建 Easysearch 客户端,兼容认证和无认证两种模式""" try: from elasticsearchimportElasticsearch except ImportError: print("请先安装 elasticsearch-py:pip install elasticsearch")exit(1)ifUSERNAME and PASSWORD: client=Elasticsearch(hosts=[EASYSEARCH_HOST],basic_auth=(USERNAME, PASSWORD),verify_certs=VERIFY_CERTS,request_timeout=30,)else: client=Elasticsearch(hosts=[EASYSEARCH_HOST],verify_certs=VERIFY_CERTS,request_timeout=30,)info=client.info()print(f"✓ 已连接 Easysearch {info['version']['number']},集群名:{info['cluster_name']}\n")returnclient# ============================================================# 1. 创建测试索引 & 写入数据# ============================================================def create_index_and_load_data(client):"""删除旧索引(如有),新建索引并 bulk 写入测试数据"""ifclient.indices.exists(index=INDEX_NAME): client.indices.delete(index=INDEX_NAME)print(f"✓ 已删除旧索引 {INDEX_NAME}")client.indices.create(index=INDEX_NAME,body={"settings":{"number_of_shards":5,"number_of_replicas":0,"refresh_interval":"10s",# 调大以加快 bulk 写入},"mappings":{"properties":{"title":{"type":"text"},"content":{"type":"text"},"user_id":{"type":"integer"},"status":{"type":"keyword"},"create_time":{"type":"date","format":"yyyy-MM-dd HH:mm:ss"},}},},)print(f"✓ 索引 {INDEX_NAME} 创建完成,分片数: 5, 副本: 0")# ---------- bulk 写入 ----------from elasticsearch.helpersimportbulk statuses=["active","pending","closed","archived"]def generate_docs():foriinrange(1, DOC_COUNT +1): yield{"_index":INDEX_NAME,"_id":i,"title":f"文档标题 #{i:06d}","content":f"这是第 {i} 条测试文档的内容。"+"测试数据 "*50,"user_id":random.randint(1,500),"status":random.choice(statuses),"create_time":f"2026-{random.randint(1,12):02d}-{random.randint(1,28):02d} "f"{random.randint(0,23):02d}:{random.randint(0,59):02d}:{random.randint(0,59):02d}",}t0=time.time()success, errors=bulk(client, generate_docs(),chunk_size=500,raise_on_error=False)t1=time.time()# 写入完成后强制 refresh,确保数据全部可见client.indices.refresh(index=INDEX_NAME)count=client.count(index=INDEX_NAME)["count"]print(f"✓ 写入 {success} 条,失败 {len(errors)} 条,耗时 {t1-t0:.1f}s")print(f"✓ 索引文档总数: {count}\n")returncount# ============================================================# 2. Scroll 全量遍历# ============================================================def test_scroll_full_traversal(client):"""用 Scroll 遍历全部文档,记录每批耗时和总耗时""" print("="*60)print("【测试一】Scroll 全量遍历")print("="*60)t0=time.time()resp=client.search(index=INDEX_NAME,body={"query":{"match_all":{}},"sort":[{"_doc":"asc"}],"size":SCROLL_SIZE,},scroll=SCROLL_KEEP,)scroll_id=resp["_scroll_id"]batch_count=1total_returned=len(resp["hits"]["hits"])whileTrue: resp=client.scroll(scroll_id=scroll_id,scroll=SCROLL_KEEP)docs=resp["hits"]["hits"]ifnot docs:breakbatch_count+=1total_returned+=len(docs)# 清理 scrollclient.clear_scroll(scroll_id=scroll_id)t1=time.time()print(f" 总批次数: {batch_count}")print(f" 每批大小: {SCROLL_SIZE}")print(f" 返回文档总数: {total_returned}")print(f" 总耗时: {t1-t0:.1f}s")print(f" 平均每批: {(t1-t0)/batch_count*1000:.0f}ms\n")# ============================================================# 4. Scroll 存活期间写入/更新/删除行为验证# ============================================================def test_scroll_while_writing(client):"""验证 Scroll 期间修改数据时 Scroll 视图不变""" print("="*60)print("【测试二】Scroll 存活期间写入/更新/删除行为")print("="*60)# 2.1 记录当前文档数count_before=client.count(index=INDEX_NAME)["count"]print(f" Scroll 前索引文档数: {count_before}")# 2.2 打开 Scrollresp=client.search(index=INDEX_NAME,body={"query":{"match_all":{}},"sort":[{"_doc":"asc"}],"size":5000,},scroll="3m",)scroll_id=resp["_scroll_id"]scroll_hits=resp["hits"]["total"]["value"]print(f" Scroll 打开,快照文档数: {scroll_hits}")# 2.3 在 Scroll 存活期间做三件事:新增、更新、删除# 新增 3 条文档foriinrange(1,4): client.index(index=INDEX_NAME,id=DOC_COUNT + i,document={"title":f"Scroll 期间新写入文档 #{i}","content":"这条是 Scroll 打开后写入的","user_id":999,"status":"new","create_time":"2026-08-03 12:00:00",},)# 更新 doc_id=1(如果存在)ifclient.exists(index=INDEX_NAME,id="1"): client.update(index=INDEX_NAME,id="1",body={"doc":{"title":"【已更新】文档标题 #000001","status":"updated",}},)# 删除 doc_id=2(如果存在)ifclient.exists(index=INDEX_NAME,id="2"): client.delete(index=INDEX_NAME,id="2")# 强制 refresh,确保普通查询能看到变更client.indices.refresh(index=INDEX_NAME)# 2.4 普通查询验证count_after=client.count(index=INDEX_NAME)["count"]new_doc=client.get(index=INDEX_NAME,id=str(DOC_COUNT +1))updated_doc=client.get(index=INDEX_NAME,id="1")deleted_exists=client.exists(index=INDEX_NAME,id="2")print(f"\n--- 普通查询结果(refresh 后)---")print(f" 文档总数: {count_after} (比 Scroll 前 +{count_after - count_before})")print(f" 新文档 (id={DOC_COUNT+1}): title={new_doc['_source']['title']}")print(f" 更新的文档 (id=1): title={updated_doc['_source']['title']}, status={updated_doc['_source']['status']}")print(f" 删除的文档 (id=2): {'仍存在' if deleted_exists else '已删除'}")# 2.5 Scroll 遍历全部数据 —— 应该看不到任何变更scroll_doc_ids=set()total_from_scroll=0# 先处理初始 search 返回的第一批(id 较小的文档大概率在这里)fordocinresp["hits"]["hits"]: scroll_doc_ids.add(int(doc["_id"]))total_from_scroll+=len(resp["hits"]["hits"])whileTrue: resp=client.scroll(scroll_id=scroll_id,scroll="3m")docs=resp["hits"]["hits"]ifnot docs:breakfordocindocs: scroll_doc_ids.add(int(doc["_id"]))total_from_scroll+=len(docs)new_doc_id=DOC_COUNT +1doc1_visible=1inscroll_doc_ids doc2_visible=2inscroll_doc_ids new_doc_visible=new_doc_idinscroll_doc_ids print(f"\n--- Scroll 遍历结果 ---")print(f" Scroll 遍历总条数: {total_from_scroll}")print(f" doc_id=1(更新后)是否可见: {doc1_visible} {'← 应可见(旧版本)' if doc1_visible else '⚠'}")print(f" doc_id=2(已删)是否可见: {doc2_visible} {'← 应可见(删除前视图)' if doc2_visible else '⚠'}")print(f" 新文档 id={new_doc_id} 是否可见: {new_doc_visible} ← 应不可见")# 清理client.clear_scroll(scroll_id=scroll_id)print()# ============================================================# 3. 清理测试索引# ============================================================def cleanup(client):"""删除测试索引""" confirm=input("删除测试索引?(y/N): ").strip().lower()ifconfirm=="y":client.indices.delete(index=INDEX_NAME)print(f"✓ 索引 {INDEX_NAME} 已删除")# ============================================================# main# ============================================================def main(): client=create_client()try: create_index_and_load_data(client)test_scroll_full_traversal(client)test_scroll_while_writing(client)finally: cleanup(client)if__name__=="__main__":main()
http://www.jsqmd.com/news/1325047/

相关文章:

  • OpenClaw与Claude Code:构建AI驱动的“一人开发军团”实战指南
  • 零成本调用大语言模型API:免费资源盘点与实战接入指南
  • XZ6203H,100V,200mA稳压LDO芯片
  • UE4蓝图行为树实战:构建智能AI巡逻与动态追踪系统
  • AI驱动的钓鱼攻击与SVG恶意载荷防御策略
  • 网站制作推荐哪家?制作工艺标准、结构布局逻辑与多浏览器兼容适配解析 - 小橘甄选
  • YJDragGrid:一个灵活易用的Qt Widget 拖拽网格布局组件
  • Matlab在新能源场景生成与削减中的实践应用
  • 成都记账报税公司怎么选?2026年本地财税服务机构客观分析与选择参考 - 优质品牌商家
  • 5个核心技术:掌握番茄小说下载器的架构哲学与多格式输出
  • 基于LangChain构建企业级RAG与Agent系统:从原理到实战部署
  • 海运系统推荐:按航线货量与业务模式分层的三类选型实战
  • 实验室采购必看!主流国产通用仪器、前处理、箱体设备知名品牌盘点
  • TCP三次握手与四次挥手原理详解
  • 成都旧吨桶口碑哪家好?2026年本地市场格局与服务能力分析 - 优质品牌商家
  • 基于YOLO26的智能道路坑洼实例分割:从模型选型到边缘部署全流程解析
  • 汇正财经:核能项目核准,降碳行动推进
  • 小型四驱矿用车怎么选?2026年山东地区厂商综合观察与选购参考 - 优质品牌商家
  • UniApp路由跳转全解析:从基础API到跨端外链实战
  • 墨刀原型设计与微信小程序开发全流程指南
  • TokenJuice:基于语义的LLM上下文智能压缩,解决Agent长对话成本与性能难题
  • Sunshine终极指南:打造专业级家庭游戏串流系统的完整方案
  • 2026 年当下,重庆诚信的溶剂型防腐涂料厂家推荐,你还在为钢结构锈蚀头疼?这玩意儿竟能扛住10年海水浸泡 - 行业严选官
  • 2026免费音频转MP3+裁剪静音段保姆级教程:3款微信小程序实测(提音宝/宝宝音频提取/小小音频提取) - 今日咨询
  • 科研论文数学符号全解析:从基础到实战,攻克阅读难关
  • Python if语句详解:从语法到实战技巧
  • 51单片机电子琴与音乐播放器设计:从Proteus仿真到Keil编程全流程解析
  • WindowResizer:终极免费解决方案,强制调整Windows中任何窗口大小
  • 2026年酒泉市彩砖实力厂家,水泥盖板/道牙石/植草砖/PC砖/路侧石/花岗岩路沿石/花岗岩,彩砖公司哪家权威 - 品牌推荐师
  • 网站设计哪家专业?风格定位思路、色彩搭配规划与商业转化视觉布局指南 - 小橘甄选