本文介绍如何在阿里云 Elasticsearch 中挂载 Apache Paimon 表,并基于 Paimon global index 对湖表数据进行向量检索、标量过滤和字段回源查询。
实现方式
Paimon 表中的原始数据仍存储在 OSS 等湖仓存储中。Elasticsearch 读取 Paimon 表快照和 global index 元信息后,将湖表索引挂载为一个可查询的 Elasticsearch 索引。业务侧可以继续使用 Elasticsearch Query DSL、KNN 查询、过滤条件和字段返回能力访问湖上数据,无需将全量原始数据同步到在线 Elasticsearch 索引中。实现方式如下:
-
将 Paimon 表中的向量列构建为 Elasticsearch 可查询的向量索引。
-
将业务 ID、标签、URI、时间、类目等标量列写入 global index,用于过滤和返回。
-
通过
/_paimon/mount将 Paimon 表挂载为 Elasticsearch 索引。 -
使用 Elasticsearch KNN 查询完成向量召回,并结合
term、terms、range等过滤条件。 -
在命中 TopK 后按需读取 Paimon 表中的字段,减少在线索引对原始数据的重复存储。
-
使用 Spark SQL
vector_search在 Paimon 侧验证同一条 query vector 的召回结果。
典型链路如下:
业务数据 / 图片 / 文本 / Embedding
-> 写入 Paimon 表
-> 构建 Paimon global index
-> Elasticsearch 挂载 Paimon 表
-> 使用 ES Query DSL / KNN 查询
-> 命中后返回标量列或按需回源字段
适用场景
该功能适用于以下场景:
-
图片、文本、视频帧等多模态数据存储在 Paimon/OSS 中,希望通过 Elasticsearch 提供在线相似检索服务。
-
业务需要同时使用向量检索和结构化过滤,例如“相似图片 + 类目过滤”“相似文档 + 标签过滤”。
-
不希望把湖表中的全量字段重复同步到 Elasticsearch,只希望保留搜索索引和必要的返回字段。
-
需要用 Elasticsearch 的查询 DSL、排序、过滤、聚合、Kibana 或 RAG 应用生态作为湖上数据的搜索入口。
前提条件
下表为使用该能力的整体要求。本文步骤一、步骤二会以示例表 paimon_vector_demo 演示从零创建 Paimon 表并构建 es-index global index;如果已有满足要求的 Paimon 表和 es-index global index,可直接从步骤三开始。
|
类型 |
要求 |
|
Elasticsearch 实例 |
实例版本需 ≥9.4,镜像已内置 Paimon mount、source 及 global index 读取能力。
|
|
Paimon 表 |
表数据已写入 OSS 或 DLF 管理的湖仓存储,且表中包含可用于检索的向量列。 |
|
Global index |
已为 Paimon 表构建 |
|
权限 |
Elasticsearch 实例或挂载请求中使用的凭据需要具备读取 Paimon 表路径、manifest、数据文件和 index 文件的权限。 |
|
Spark 环境 |
如需构建索引或使用 Spark SQL 验证,需要可运行 Paimon Spark 扩展和对应的 Paimon/ES 集成 Jar。 |
使用限制和注意事项
-
如果需要通过
_source或 Paimon source 回填字段,建表和写入链路需要正确开启row-tracking.enabled=true和data-evolution.enabled=true。否则数据文件可能缺少firstRowId,导致回源阶段无法建立 rowId 到数据文件的映射。 -
需要参与过滤或写入 global index 的伴随标量列,必须在 Spark SQL 的
index_column中与向量列一起声明。只把字段放入return_fields或_source,不能保证其可用于term、range等过滤查询。 -
如果向量字段使用
dot_product作为相似度,写入表中的向量和查询向量通常都需要先做 L2 normalize,确保分数和排序可比。 -
HNSW、DiskBBQ 等向量索引算法的维度上限、召回参数和性能特征以当前 Elasticsearch 实例版本为准。
-
大规模 HNSW 索引如果使用
heap挂载模式,需要提前评估 JVM heap、shard 数和向量图大小。功能验证建议优先使用mmap或hybrid模式。 -
Spark SQL
vector_search的查询向量需要写成 Spark 可解析的 Float/Double 数组字面量,例如array(1.234E-2D, -3.456E-2D),避免被解析为 DECIMAL 或表达式。
步骤一:准备 Paimon 表
以下示例以 OSS filesystem catalog 为例。DLF catalog 场景下,请将 catalog 和表路径替换为您的实际配置。
配置 Spark Catalog
# 使用 Paimon Spark 扩展
spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions
# 配置 OSS filesystem catalog
spark.sql.catalog.oss=org.apache.paimon.spark.SparkCatalog
spark.sql.catalog.oss.metastore=filesystem
spark.sql.catalog.oss.warehouse=oss://<bucket>/<warehouse>
spark.sql.catalog.oss.fs.oss.endpoint=oss-<region>-internal.aliyuncs.com
spark.sql.catalog.oss.fs.oss.accessKeyId=${OSS_AK_ID}
spark.sql.catalog.oss.fs.oss.accessKeySecret=${OSS_AK_SECRET}
# Hadoop OSS 读写配置
spark.hadoop.fs.oss.endpoint=oss-<region>-internal.aliyuncs.com
spark.hadoop.fs.oss.accessKeyId=${OSS_AK_ID}
spark.hadoop.fs.oss.accessKeySecret=${OSS_AK_SECRET}
# 避免 Paimon v1 function 抢占外部函数解析
spark.paimon.v1Function.enabled=false
创建示例表
以下示例表包含:
-
id:业务主键。 -
label:用于标量过滤的标签列。 -
content_uri:可返回给业务的对象路径。 -
emb_norm:归一化后的向量列。
CREATE TABLE oss.default.paimon_vector_demo (
id BIGINT,
label BIGINT,
content_uri STRING,
emb ARRAY<DOUBLE>,
emb_norm ARRAY<FLOAT>
) TBLPROPERTIES (
'write-mode' = 'append-only',
'bucket' = '-1',
'file.format' = 'parquet',
'row-tracking.enabled' = 'true',
'data-evolution.enabled' = 'true',
'global-index.es-index.fields.emb_norm.algorithm' = 'hnsw',
'global-index.es-index.fields.emb_norm.dimension' = '768',
'global-index.es-index.fields.emb_norm.metric' = 'dot_product'
);
如果使用图片、文本等多模态数据,建议先在离线链路中生成 embedding,再写入 Paimon 表。使用dot_product时,建议写入归一化后的向量列,例如emb_norm。
写入数据并检查
将源数据写入示例表,并对向量列做 L2 归一化。以下 source_table 为您自有的源数据表(包含 id、label、content_uri、emb 等列),请替换为实际来源:
INSERT INTO oss.default.paimon_vector_demo
SELECT
id,
label,
content_uri,
emb,
vector_l2_normalize(emb) AS emb_norm
FROM source_table;
vector_l2_normalize依赖对应版本的 Paimon Spark 扩展。按约定删除SQL中函数vector_l2_normalize,要求上游生成归一化后的ARRAY<FLOAT>,纯 SQL 写法可保留。
INSERT INTO oss.default.paimon_vector_demo
SELECT
id,
label,
content_uri,
emb,
transform(
emb,
x -> CAST(x / sqrt(aggregate(emb, CAST(0 AS DOUBLE), (acc, v) -> acc + v * v)) AS FLOAT)
) AS emb_norm
FROM source_table;
写入后检查行数、维度和向量归一化结果:
SELECT count(*) AS total_count
FROM oss.default.paimon_vector_demo;
SELECT
id,
size(emb_norm) AS dim,
aggregate(emb_norm, CAST(0 AS FLOAT), (acc, v) -> acc + v * v) AS sumsq
FROM oss.default.paimon_vector_demo
LIMIT 10;
归一化正确时,dim 应等于建表设置的维度(如 768),sumsq 应接近 1.0。
步骤二:构建 Paimon global index
推荐在 EMR Serverless Spark SQL 中调用 Paimon 的 create_global_index 存储过程构建 es-index(示例中为 oss.sys.create_global_index,其中 OSS 为步骤一配置的 catalog 名)。后续新版本 EMR Serverless Spark 会更新并内置支持该能力的最新 Paimon;存量老版本不会更新内置 Paimon,如需使用该能力,必须手动打包兼容的 Paimon、paimon-eslib 及 ESLib/Lucene 依赖,并通过作业配置排除和替换内置 Paimon:可使用官方提供的 emr-3.5 适配 Jar,或用社区最新版本自行打包。当前托管 DLF 自动构建产物的索引类型与 Elasticsearch paimon-store 所需的 es-index 不一致,暂不能直接替代本步骤;DLF catalog 仍可用于管理 Paimon 表元数据。
# 老版本 EMR Serverless Spark:排除并替换内置 Paimon
spark.emr.serverless.excludedModules paimon
spark.emr.serverless.user.defined.jars oss://<bucket>/jars/<paimon-es-index-bundle>.jar
spark.sql.extensions org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions
# 使用 Spark SQL 构建 es-index
CALL oss.sys.create_global_index(
table => 'default.paimon_vector_demo',
index_column => 'emb_norm,id,label,content_uri',
index_type => 'es-index',
options => 'global-index.row-count-per-shard=100000,global-index.es-index.fields.emb_norm.m=30,global-index.es-index.fields.emb_norm.ef_construction=360'
);
Spark 配置和 SQL 参数说明如下。
|
参数 |
说明 |
|
|
esr-5.3.1 及以下需设置为 paimon,排除 EMR Serverless Spark 内置的 Paimon 模块,避免新旧类冲突。 |
|
|
自定义 JAR 的 OSS 地址。需包含与 Spark 版本匹配的 Paimon、paimon-eslib 及 ESLib/Lucene 依赖:可使用官方提供的 emr-3.5 适配 Jar,或用社区最新版本自行打包。 |
|
|
必须配置为 org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions。 |
|
|
目标 Paimon 表名。调用 oss.sys procedure 时填写 default.paimon_vector_demo。 |
|
|
写入同一个 es-index 的列。主向量列 emb_norm 必须放第一位,后续为过滤或返回的伴随标量列。 |
|
|
固定使用 es-index,生成当前 Elasticsearch paimon-store 可识别的索引文件。 |
|
|
传入分片行数及 HNSW m、ef_construction 等构建参数;调用参数优先于表属性。 |
|
重复构建 |
首次调用构建现有数据;后续调用补建未覆盖的新数据。如需强制全量重建,先 drop_global_index,再重新调用 create_global_index。 |
构建完成后,可通过 Paimon 系统表检查 global index 文件:
SELECT
index_type,
index_field_name,
row_count,
file_name
FROM oss.default.`paimon_vector_demo$table_indexes`
WHERE index_type = 'es-index';
paimon_vector_demo$table_indexes系统表的可用列为file_name、file_size、bucket、index_type、row_count。请使用file_name查看索引文件,而非file_path。
步骤三:将 Paimon 表挂载为 Elasticsearch 索引
调用 /_paimon/mount 接口,将 Paimon 表挂载为 Elasticsearch 索引。
POST /_paimon/mount
{
"auth_type": "direct",
"table_path": "oss://<bucket>/<warehouse>/default.db/paimon_vector_demo",
"oss_endpoint": "oss-<region>-internal.aliyuncs.com",
"oss_bucket": "<bucket>",
"oss_access_key_id": "${OSS_AK_ID}",
"oss_access_key_secret": "${OSS_AK_SECRET}",
"index_name": "paimon_vector_demo_index",
"vector_field_name": "emb_norm",
"storage_mode": "mmap",
"source_enabled": true,
"return_fields": [
"id",
"label",
"content_uri"
],
"timeout": "600"
}
调用成功后返回类似如下结果:
{
"acknowledged": true,
"alias": "paimon_vector_demo_index",
"index": "paimon_vector_demo_index_2"
}
mount 会创建一个实际索引(名称为index_name加数字后缀,例如paimon_vector_demo_index_2),并将别名index_name指向该索引。后续查询使用别名paimon_vector_demo_index即可;步骤四中_cat/indices显示的是带数字后缀的实际索引名。
参数说明如下。其中 table_path、oss_endpoint、oss_bucket、oss_access_key_id、oss_access_key_secret 为必填参数,缺失时接口会返回 Required fields 错误。
|
参数 |
是否必填 |
说明 |
|
|
否 |
认证方式。示例使用 direct,即在请求中传入 OSS 访问凭据。生产环境建议优先使用更安全的托管凭据或 RAM 角色方式。 |
|
|
是 |
Paimon 表路径。 |
|
|
是 |
OSS endpoint。建议使用与 Elasticsearch 实例同地域的内网 endpoint。 |
|
|
是 |
表所在 OSS bucket。 |
|
|
是 |
访问 OSS 的 AccessKey ID。 |
|
|
是 |
访问 OSS 的 AccessKey Secret。 |
|
|
否 |
挂载后生成的 Elasticsearch 索引名称。 |
|
|
否 |
用于向量检索的字段名称。 |
|
|
否 |
索引读取模式,例如 mmap、hybrid 或 heap。可用值以当前版本为准。 |
|
|
否 |
是否启用 Paimon source 回源。 |
|
|
否 |
需要从 Paimon 表中返回的字段。 |
|
|
否 |
mount 超时时间,单位为秒。 |
建议首次验证使用mmap或hybrid。如果需要使用heap模式做性能压测,请先评估索引规模和 JVM heap。
步骤四:检查挂载结果
挂载完成后,检查索引状态、文档数量、mapping 和 settings。
GET /_cat/indices/paimon_vector_demo_index?v
GET /paimon_vector_demo_index/_count
GET /paimon_vector_demo_index/_mapping
GET /paimon_vector_demo_index/_settings?flat_settings=true
重点检查以下内容:
-
索引状态为
green或符合预期。 -
_count与 Paimon 表行数一致。 -
向量字段 mapping 中维度和相似度配置正确,例如
dims=768、similarity=dot_product。 -
标量字段
id、label、content_uri已出现在 mapping 中。 -
settings 中 Paimon 表路径、snapshot ID、storage mode 等信息符合预期。
步骤五:使用 Elasticsearch 查询
纯向量查询
POST /paimon_vector_demo_index/_search
{
"size": 10,
"_source": [
"id",
"label",
"content_uri"
],
"knn": {
"field": "emb_norm",
"query_vector": [/* 768 维归一化 query vector */],
"k": 10,
"num_candidates": 100
}
}
说明:挂载索引的标量列(id、label、content_uri)通过_source回源返回(对应 mount 时的return_fields),请使用_source指定需要返回的字段。这些列不是 Lucene stored field/doc_values,使用"_source": false搭配fields检索取不到值(命中项仅含_index和_score)。如需连同向量列一起返回,可将_source设为true。
向量 + 标量过滤
POST /paimon_vector_demo_index/_search
{
"size": 10,
"_source": [
"id",
"label",
"content_uri"
],
"knn": {
"field": "emb_norm",
"query_vector": [/* 768 维归一化 query vector */],
"k": 10,
"num_candidates": 100,
"filter": {
"term": {
"label": 1
}
}
}
}
使用 terms 或 range 过滤
将上一个请求中的 knn.filter 替换为以下内容即可。
{
"terms": {
"label": [1, 2, 3]
}
}{
"range": {
"id": {
"lt": 100000
}
}
}
返回 Paimon 表字段
如果已启用 source_enabled,可返回 Paimon 表中的字段:
POST /paimon_vector_demo_index/_search
{
"size": 10,
"_source": true,
"knn": {
"field": "emb_norm",
"query_vector": [/* 768 维归一化 query vector */],
"k": 10,
"num_candidates": 100
}
}
如果 _source 为空,请重点检查:
-
建表时是否开启
row-tracking.enabled=true和data-evolution.enabled=true。 -
写入路径是否正确分配数据文件的
firstRowId。 -
return_fields中是否包含需要返回的字段。 -
Elasticsearch 实例是否已加载 Paimon source 能力。
步骤六:使用 Spark SQL 验证查询结果
您也可以在 Spark SQL 中使用同一条 query vector 验证 Paimon global index 查询结果。
USE oss.default;
SELECT
id,
label,
content_uri,
__paimon_search_score AS score
FROM vector_search(
'oss.default.paimon_vector_demo',
'emb_norm',
array(/* 768 double literals, e.g. 1.234567E-2D */),
10,
map('hnsw.num_candidates', '100')
)
ORDER BY score DESC
LIMIT 10;
向量 + 标量过滤示例:
SELECT
id,
label,
content_uri,
__paimon_search_score AS score
FROM vector_search(
'oss.default.paimon_vector_demo',
'emb_norm',
array(/* same normalized query vector */),
10,
map('hnsw.num_candidates', '100')
)
WHERE label = 1
ORDER BY score DESC
LIMIT 10;
使用 EXPLAIN FORMATTED 检查是否走向量索引:
EXPLAIN FORMATTED
SELECT
id,
label,
__paimon_search_score AS score
FROM vector_search(
'oss.default.paimon_vector_demo',
'emb_norm',
array(/* same normalized query vector */),
10,
map('hnsw.num_candidates', '100')
)
WHERE label = 1
ORDER BY score DESC;
预期计划中出现 vector search、global index scan 或 VectorSearch 相关信息。
常见问题
为什么 KNN 查询能返回结果,但标量过滤查不到?
请检查过滤字段是否已加入 Spark SQL 的 index_column。例如需要按 label 过滤时,构建时必须包含 index_column => 'emb_norm,id,label,content_uri'。如果只在 return_fields 中配置 label,该字段可以作为回源返回字段,但不一定能作为索引过滤字段。
为什么 _source 返回为空?
常见原因是 Paimon 数据文件缺少 rowId 到文件的映射信息。请确认:
-
建表时开启
row-tracking.enabled=true。 -
建表时开启
data-evolution.enabled=true。 -
写入链路使用支持 row tracking 的路径。
-
旧表或旧数据文件可能需要重新写入或重新构建。
查询向量应该使用原始向量还是归一化向量?
如果索引字段为归一化向量列 emb_norm,且相似度为 dot_product,查询向量也应使用同模型生成并归一化后的向量。否则可能出现分数不可比、排序不符合预期,甚至被向量维度或取值校验拒绝。
mmap、hybrid 和 heap 应该怎么选?
-
mmap:适合功能验证和较低内存占用场景,依赖操作系统 page cache。 -
hybrid:在内存与磁盘访问之间折中,适合中大规模索引的常规验证。 -
heap:适合性能专项,但会占用较多 JVM heap。使用前需要根据向量数量、维度、shard 数和副本数评估内存。