Lake Search:ES x Paimon 让湖上多模态数据可搜可用

简介: 当图片、视频、文本和向量在 Paimon 中增长到 PB 级,传统“同步到湖外再建索引”的方式,会让搜索面临数据就绪慢、第二份事实数据成本高和版本治理复杂等问题。本文介绍阿里云 Elasticsearch 9.4 Search Lake 如何直接挂载与 Paimon 表版本关联的 Global Index,在不复制事实数据的前提下提供 BM25、kNN、结构化过滤、排序与聚合能力,并结合多模态样本湖场景拆解方案架构、Demo 及性能与成本取舍。关键词:Search Lake、Apache Paimon、Elasticsearch 9.4、Global Index、OpenLake、多模态检索

当图片、视频、文本和向量在 Paimon 中增长到 PB 级,传统“同步到湖外再建索引”的方式,会让搜索面临数据就绪慢、第二份事实数据成本高和版本治理复杂等问题。本文介绍阿里云 Elasticsearch 9.4 Search Lake 如何直接挂载与 Paimon 表版本关联的 Global Index,在不复制事实数据的前提下提供 BM25、kNN、结构化过滤、排序与聚合能力,并结合多模态样本湖场景拆解方案架构、Demo 及性能与成本取舍。

关键词:Lake Search、Apache Paimon、Elasticsearch 9.4、Global Index、OpenLake、多模态检索

1. 数据入湖,索引如何“出湖”?

1.1 PB 级数据入湖后,如何进入低延迟搜索?

Apache Paimon 正从流式湖表格式演进为面向 Data + AI 的统一多模态数据湖底座:Paimon 1.0 已在生产环境验证超过 100PB 的存储规模,淘宝、天猫每天还有近 10PB 的图像、音频和视频等多模态数据写入 Paimon。围绕这些数据,Flink 负责 CDC 和实时入湖,Spark 负责批量加工、特征和向量生成,StarRocks 等引擎提供交互式 OLAP 查询,湖仓体系已经解决了数据存储、治理和计算的问题。

问题出现在这些数据要被在线应用“找到”的时候。以具身智能训练样本湖为例,一次模型评测暴露了“机械臂抓取透明物体后滑落”的问题。下一轮训练开始前,算法工程师需要从长期积累的视频中找到相似失败片段,同时检索操作指令与人工标注,按模型版本、时间和复核状态过滤,并查看失败原因分布。一次请求需要组合视频向量召回、BM25 全文相关性、结构化过滤、排序、聚合和 RRF 融合,还要在交互过程中保持稳定低延迟。

真正的瓶颈并不是某条 SQL 能不能写,也不是能不能完成一次向量 TopK,而是这条搜索链路能否在 PB 级数据上持续服务在线应用。具体表现为三个问题:

  • 数据就绪与查询时延:最新样本写入 Paimon 后,如果还要全量读取、跨系统传输并重建索引,搜索侧会长期落后于湖表版本;数据越大,等待越长,也越难同时满足近实时就绪、低延迟和高 QPS。
  • 第二份事实数据与同步成本:湖表和搜索系统各自保存一份业务数据后,不仅增加存储,还要长期维护增量同步、版本一致性、权限对齐、失败重试和历史回补。
  • 搜索游离于 OpenLake 多引擎体系之外:Spark、Flink 和 StarRocks 可以围绕同一份湖表完成加工与分析,搜索却仍依赖湖外副本,应用需要在两套数据版本和治理链路之间反复对齐。

这三个问题共同指向同一矛盾:数据已经进入湖,在线搜索仍依赖数据再次“出湖”。Search Lake 要解决的,是让 ES 作为 OpenLake 的搜索引擎直接使用与 Paimon 表版本关联的 Global Index;具体如何实现,后文展开。

image.png

图 1:一次多模态样本检索背后的湖上搜索链路。


1.2 湖内分析与在线搜索:不同访问路径的取舍

面对这类需求,现有方案通常沿着两条路径展开:直接使用湖内分析引擎查询,或者把数据同步到湖外搜索系统。两条路径都能解决一部分问题,但各有明确的执行目标和代价。

湖内分析路径擅长 SQL、Join、分组统计、报表和 ETL,强调在较大数据范围内得到完整、可复现的计算结果。分析引擎也可以完成过滤、聚合,甚至提供向量函数,但相关性排序、多路召回融合、连续筛选和高并发在线查询并不是这条路径的主要执行模型。

湖外搜索路径可以提供 BM25、kNN、过滤、排序和聚合等在线检索能力,但传统接入方式通常需要从湖表读取数据,经 ETL 或流任务同步到 ES、向量库等系统,再重新构建和维护索引。搜索能力有了,湖外第二份数据及其生命周期也随之产生。

因此,真正需要比较的不是 SQL 和搜索谁能写出某条查询,而是两种系统如何分工:计算引擎负责把数据算全、加工好,搜索引擎负责快速返回最相关的一批结果并承接在线流量。

Search Lake 选择把搜索索引留在湖上,并让 ES 直接挂载与 Paimon 表版本对应的 Global Index。图 2 对比了传统湖外同步和湖上索引挂载两条链路;

image.png

图 2:从湖外数据副本到湖上搜索索引。

PB 级数据就绪:从全量复制到挂载湖上索引

当数据规模进入 PB 级,两条链路的差异会首先体现在“多久可以开始查询”。为了把量级差异转成可核查的模型,下面选取 1 PB、2 PB 和 5 PB 三档数据量:传统离线导入分别按 10 GB/s 理想吞吐,以及 4 GB/s 有效吞吐并计入 20% refresh、校验开销估算;Search Lake mount 按索引文件数、分片数和元数据打开并发,给出 12、22 和 45 分钟的方案量级。

图 3:1 PB、2 PB 和 5 PB 数据规模下,传统离线导入与 Search Lake mount 的数据就绪时间估算。


2. Search Lake 是什么?

Search Lake 是 Paimon 与 Elasticsearch 共同提供的湖上搜索方案。从阿里云 Elasticsearch 9.4 开始,客户可以将与 Paimon 表版本关联的 Global Index 直接挂载为 ES 索引,并通过熟悉的查询接口获得全文、向量、过滤、聚合、排序和按需回源能力。

它的边界可以用三句话说明:

  • 数据仍以 Paimon 表版本为准。Paimon 的底层存储可以是 OSS、S3 或其他对象存储,DLF 继续管理目录、权限和血缘。
  • 搜索索引仍在湖上。Global Index 随表版本记录索引文件、分片范围和行号关系,不需要为 ES 维护一份新的事实数据副本。
  • ES 负责把湖上索引变成应用可用的搜索接口。应用继续使用 Query DSL、KNN 和 aggregation,命中后再按需读取业务字段和原始对象

通过 ES 和 Paimon 的结合,可以为 PB 规模以上的多模态数据湖提供在线、离线一体化的检索方案,兼顾成本、性能和效果。


2.1 为什么 ES 适合作为湖上搜索入口?

ES 适合作为湖上搜索入口,首先是因为它把“如何找到最相关的一批结果”作为核心执行模型。它不是用来替代 Spark 的扫描和计算,而是把搜索需要的索引、相关性、TopK 和交互式下钻组织在同一个引擎中。

  • 倒排索引和 BM25 用于全文相关性检索。
  • doc values、bitset 和 query cache 用于结构化过滤、排序与聚合。
  • HNSW、BBQ 和 DiskBBQ 用于不同规模与成本目标下的向量召回。
  • 分片并行、TopK 剪枝和 segment 级执行用于降低在线查询延迟。
  • Query DSL、ES|QL、Kibana 与 RAG / AI Search 生态让搜索结果更容易进入应用。

ES 将交互式查询压到百毫秒级

为了验证这些执行机制在湖上数据中的效果,我们在同一份 Cohere 10M 数据上设计了两组语义等价的结构化查询。Spark SQL 显式关闭 Global Index,ES 查询挂载后的只读索引;两侧使用相同过滤条件、排序和返回行数,每组执行 1 次预热和 5 次计量,并对返回行数和规范化结果 SHA-256 逐项校验。

过滤、排序并返回 Top100 时,Spark SQL p50 为 549.433 ms,ES p50 为 208.908 ms,相差 2.63 倍;范围过滤并分组聚合时,Spark SQL p50 为 1644.415 ms,ES p50 为 76.076 ms,相差 21.62 倍。

image.png

图 4:等价结果下,ES 与 Spark SQL 的交互式查询时延。

面向对象存储的向量索引路径

Paimon 通常把表数据和 Global Index 持久化在 OSS、S3 等对象存储中。若没有面向搜索的索引路径,一次向量 TopK 需要定位并读取大量向量数据,容易受到对象打开、随机读取、反序列化、读放大和缓存未命中的影响,性能退化到秒级/分钟级。

DiskBBQ 解决的是索引阶段的访问效率:它先用压缩表示和聚类信息缩小候选范围,把高频元信息放在内存或本地缓存中,把体量较大的候选列表和向量块留在对象存储;查询只读取少量相关块并进行重排。这样,对象存储上的冷数据不必全部变成热副本,也能获得可调节的延迟、召回与内存成本。

image.png

图 5:DiskBBQ 面向对象存储的向量检索路径。


2.2 技术架构:Paimon Global Index 与 ES 索引系统对接

从技术上看,这不是三次顺序执行的导入任务,而是 Paimon 与 ES 在三处建立同一套索引视图:构建阶段对齐分片和行号,挂载阶段把指定表版本转换为只读 ES 索引,查询阶段再用同一行号按需读取事实数据。下面分别展开。

2.2.1 索引构建与湖上持久化:Paimon Global Index + ESLib

image.png

图 6:Paimon Global Index 在湖内并行构建 ES 原生索引。

索引构建由 Spark 或 Flink 任务并行完成。Paimon Global Index 负责定义索引字段、划分分片与数据范围;每个构建任务加载 ESLib,将向量、文本和结构化字段写入同一条 Lucene 文档,生成 ES 可直接读取的 Lucene 分片。

构建阶段同时完成两项对齐:Global Index 分片与 ES 分片一一对应,分片内的 Paimon 相对 row id 与 Lucene doc id 对齐。ES 命中结果后可以直接定位对应的 Paimon 行,无需额外维护映射表。

构建完成后,Lucene 索引文件和 Global Index 元数据统一写回对象存储,并与 Paimon 表版本关联。搜索索引因此仍是湖上资产,后续 mount 可以直接复用,无需 bulk 导入、重新构建索引或复制事实数据。

2.2.2 ES 如何直接打开湖上索引

image.png

图 7:ES mount 根据湖表元数据建立只读 shard 视图,并按需读取 OSS 上的 Lucene 文件。

调用 mount 后,ES 读取 Paimon 表元数据和 Global Index 元数据,生成 mapping,并将 Global shard 映射为只读 ES shard。Lucene 文件继续保存在 OSS,无需 bulk 导入或复制索引。

ES 根据结构目录直接定位倒排索引、doc values 和向量索引等文件,无需扫描或下载完整 archive。常用元数据和 DiskBBQ 质心保留在内存,已访问的索引块进入本地缓存,其余内容按需从 OSS 读取。

节点级 buffer 上限、缓存块复用、文件页跳过、Range 合并和异步预取用于减少远端 I/O。新版本就绪后可以切换稳定索引名,应用无需感知底层变化。

用户查询实测:从全量扫描到索引服务

同一份 Cohere 10M 数据与索引保存在 OSS,查询 TopK=10。PyPaimon 精确扫描的用户观测时延为 129.550 s;Lumina fresh reader 为 7.887 s;Lumina list64 在 200 条查询批量执行时摊销为 3.960 s/query,recall@10 为 94.80%;ES mount 使用 DiskBBQ、visit=9% 时 p50 为 1.341 s,recall@10 为 95.10%。

image.png

图 8:同一份 Cohere 10M 湖上数据的用户查询路径。

这不是“Python 与 Java”或“Lumina 与 DiskBBQ 内核”的直接倍数比较,而是用户从不同入口发起查询时实际经历的完整路径。精确扫描给出无索引下限,Lumina 展示湖内 SDK 索引查询,ES mount 展示索引被挂载为在线搜索服务后的端到端请求;

2.2.3 先搜索,再按需回源

image.png

图 9:ES 先在索引中收敛 TopK,再按命中范围读取 Paimon 原始记录。

ES 先在索引中完成检索,只保留 TopK 命中,不扫描 Paimon 数据文件。

每条命中通过 shard id + doc id 直接还原为 Paimon row id,无需维护额外映射表。ES 汇总命中并合并相邻 row id,只打开相关文件和数据页,再按返回字段读取所需列。

因此,回源成本主要取决于 TopK 的数量、分布和返回字段,而不是整表规模。搜索与回源共享 Paimon 表元数据和文件目录,避免重复读取和全表扫描。


3. 使用场景与 Demo:多模态样本湖如何落地

第 1.1 节的多模态样本检索只是查询入口。真正落地时,需要先把分散在视频片段、关键帧、图片、文本说明、标注记录和业务标签中的信息组织成一张可搜索的湖表。

在 Paimon 中,每条样本记录保留样本 ID、来源、时间、场景标签、标注状态、模型版本、权限范围和对象 URI 等业务字段;视频帧、图片、文本描述和说明文档生成对应的 embedding 或全文字段。原始对象仍保存在表所使用的对象存储中,Global Index 只保存搜索需要的索引结构和回源关系。

构建索引时,向量字段负责支持“找相似视频帧 / 图片”,底层可选择 HNSW、DiskBBQ 等索引实现;文本字段负责支持关键词和语义描述检索;来源、时间、标签、状态和权限等结构化字段则用于过滤、排序和聚合分析。这些字段在同一条样本记录上对齐,业务查询不再分别调用向量库、全文搜索和湖表 SQL。

搜索结果可以直接进入下游工作流:相似样本用于训练数据筛选或内容复用,聚合分布用于发现某类场景、来源或标签的集中趋势,原图、视频片段和标注详情按需读取供人工复核。售后故障视频、自动驾驶 Corner Case、具身智能训练样本和媒资素材检索,都可以落到这条多模态样本湖链路中。

阿里云 Elasticsearch 9.4 当前支持只读挂载;Global Index 由 Spark 或 Flink 构建。搜索结果回写、标签修正闭环属于后续演进方向。



图 10:客户从 Paimon 湖表到 ES 9.4 在线查询的使用流程。


前文场景展示 Search Lake 能承接的完整业务形态;下面的 Demo 使用一组小规模售后图片数据,验证建表、构建 Global Index、挂载、组合查询和回源这条核心链路。它将这条链路压缩成一个可复现的最小闭环,操作分为五步:

  • 准备湖表:在 Paimon 中保存样本、文本、向量、标签和对象地址。
  • 构建 Global Index:通过 Spark 或 Flink 为需要搜索的向量和结构化字段建立索引。
  • 挂载到 ES 9.4:配置 DLF 与对象存储访问权限,将 Global Index 挂载为可查询的 ES 索引,无需 bulk 导入。
  • 接入搜索应用:使用 Query DSL 组合 kNN、结构化过滤和聚合。
  • 按需读取详情:根据命中结果返回业务字段、图片地址和原始对象。

视频:Paimon Global Index 构建、ES 挂载与查询演示。


4. 未来展望:从湖上搜索走向数据闭环

数据入湖以后,保存、加工和治理已经有了成熟链路,Search Lake 当前补齐的是从湖上索引到在线查询的访问路径:数据和版本继续由 Paimon 管理,索引随湖表版本保存在底层对象存储中,ES 直接挂载索引,通过全文检索、向量召回、结构化过滤和聚合分析让湖上数据进入在线应用。这条链路省去了导出、同步、重新建库,以及长期维护第二份业务数据副本的成本。

在此基础上,Search Lake 后续将从“可搜索”继续走向“可交互”。在多模态样本湖等场景中,用户检索视频和图片、复核结果、修正标签或标记误召回后,这些反馈可以重新沉淀到 Paimon,形成新的数据版本,并进入下一轮训练、评测或业务处理。一次搜索不再只是返回结果,也可以成为数据持续完善的一环。

当查询从团队内部的数据探索走向在线应用,ES 可以通过独立扩展查询节点、副本和缓存承接更高 QPS;当数据规模继续增长,对象存储、分层缓存和面向磁盘的索引能力可以减少常驻内存和重复副本。Paimon Global Index、索引构建、版本管理与 ES 查询能力也将进一步原生化,让查询性能按业务需求扩展,而不是让成本随着数据量和同步链路同步增长。从这个意义上说,索引“出湖”不是把索引搬离数据湖,而是让湖上数据进入检索、使用和反馈回湖的完整链路。


相关资源

相关文章
|
1天前
|
存储 消息中间件 人工智能
阿里云 Elasticsearch 日志采集与加工服务:让日志链路少一串组件,多一份稳定
阿里云 Elasticsearch 新版本推出日志采集与加工服务,将多源接入、流量缓冲、数据加工和可靠投递整合为云上托管能力,并与已有的写入优化、低成本存储和高性能查询能力形成完整链路,让海量日志处理变得更简单、更完整。 关键字: 阿里云 Elasticsearch、日志服务、日志采集与加工服务、读写分离、存算分离、并发查询、AI Agent
|
4月前
|
JSON 运维 Java
Apache Flink Agents 0.2.1 发布公告
Apache Flink Agents 0.2.1发布!修复3个关键缺陷(含MCP连接与Jackson反序列化问题),优化事件日志JSON输出、减小wheel包体积,并增强CI可观测性。推荐所有用户升级。支持OpenAI、Anthropic等多模型集成,附Demo演示智能运维能力。(239字)
405 5
Apache Flink Agents 0.2.1 发布公告
|
1月前
|
存储 人工智能 运维
千亿级 AI 搜索的效能实战:从混合检索到 Agentic RAG 的三年实战
本文为2026 Elastic中国大会演讲实录,直击千亿级AI搜索三大挑战:搜索融合(关键词+向量+稀疏检索原生一体)、极致效能(冷热分层、硬件降级、自研FalconSeek引擎)与Agentic RAG演进(结构化知识图谱+智能体自主推理),揭示企业级AI搜索从“能用”到“好用”再到“自进化”的实战路径。
542 8
|
1月前
|
人工智能 运维 搜索推荐
重构搜索范式:阿里云 Elasticsearch 开启“Agent 原生”时代,打造企业级 AI 记忆湖
阿里云Elasticsearch提出“Agent原生搜索”理念,打造面向AI智能体的高性能、全模态企业级AI搜索基础设施。通过Agent Skills、统一Builder平台、上下文引擎与自研FalconSeek引擎,实现结构化结果输出、分钟级Agent开发、混合检索加速及50%-300%性能提升,助力构建企业“Agent知识记忆湖”。
345 3
|
8月前
|
消息中间件 存储 Kafka
流、表与“二元性”的幻象
本文探讨流与表的“二元性”本质,指出实现该特性需具备主键、变更日志语义和物化能力。强调Kafka与Iceberg因缺乏更新语义和主键支持,无法真正实现二元性,唯有统一系统如Flink、Paimon或Fluss才能无缝融合流与表。
511 7
流、表与“二元性”的幻象
|
3月前
|
人工智能 架构师 Apache
相约深圳,全球征集|Flink Forward Asia 2026 演讲议题征集正式启动
Flink Forward Asia 2026将于6月26–27日首次落地深圳,聚焦实时计算与AI深度融合。面向全球征集议题(截止5月29日),涵盖实时AI、AI Agent、湖流一体等前沿方向。免费报名开启,共探下一代实时计算范式!
815 4
相约深圳,全球征集|Flink Forward Asia 2026 演讲议题征集正式启动
|
7月前
|
人工智能 数据处理 Apache
Forrester发布流式数据平台报告:Flink 创始团队跻身领导者行列,实时AI能力获权威认可
Ververica,由Apache Flink创始团队创立、阿里云旗下企业,首次入选Forrester 2025流式数据平台领导者象限,凭借在实时AI与流处理领域的技术创新及全场景部署能力获高度认可,成为全球企业构建实时数据基础设施的核心选择。
546 10
Forrester发布流式数据平台报告:Flink 创始团队跻身领导者行列,实时AI能力获权威认可
|
7月前
|
消息中间件 Java Kafka
在 OpenAI 打造流处理平台:超大规模实时计算的实践与思考
本文介绍OpenAI构建流处理平台的实践与挑战。面对Kafka高可用、Python生态兼容、云环境限制等问题,团队基于PyFlink打造跨区域流处理架构,集成Kafka HA组、自研代理与控制平面,支撑实时Embedding生成、特征计算等场景,并推动开源协作与平台自动化演进。
465 1
在 OpenAI 打造流处理平台:超大规模实时计算的实践与思考

热门文章

最新文章