阿里云 Elasticsearch 加速 Paimon 多模态数据湖检索实践

更新时间:
复制 MD 格式

本文介绍如何在阿里云 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 查询完成向量召回,并结合 termtermsrange 等过滤条件。

  • 在命中 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 读取能力。

  • OpenStore 配置:若镜像含 OpenStore 插件,必须在所有节点静态配置 cluster.apack.openstore.ruleout.enable: false 并重启集群(不可用动态设置替代);若无该插件则无需配置。

  • Mount 验证:执行 Mount 后,必须检查物理索引 settings,确认 index.store.type=paimon 后方可进行查询验证。

Paimon 表

表数据已写入 OSS 或 DLF 管理的湖仓存储,且表中包含可用于检索的向量列。

Global index

已为 Paimon 表构建 es-index 类型 global index。

权限

Elasticsearch 实例或挂载请求中使用的凭据需要具备读取 Paimon 表路径、manifest、数据文件和 index 文件的权限。

Spark 环境

如需构建索引或使用 Spark SQL 验证,需要可运行 Paimon Spark 扩展和对应的 Paimon/ES 集成 Jar。

使用限制和注意事项

  • 如果需要通过 _source 或 Paimon source 回填字段,建表和写入链路需要正确开启 row-tracking.enabled=truedata-evolution.enabled=true。否则数据文件可能缺少 firstRowId,导致回源阶段无法建立 rowId 到数据文件的映射。

  • 需要参与过滤或写入 global index 的伴随标量列,必须在 Spark SQL 的 index_column 中与向量列一起声明。只把字段放入 return_fields_source,不能保证其可用于 termrange 等过滤查询。

  • 如果向量字段使用 dot_product 作为相似度,写入表中的向量和查询向量通常都需要先做 L2 normalize,确保分数和排序可比。

  • HNSW、DiskBBQ 等向量索引算法的维度上限、召回参数和性能特征以当前 Elasticsearch 实例版本为准。

  • 大规模 HNSW 索引如果使用 heap 挂载模式,需要提前评估 JVM heap、shard 数和向量图大小。功能验证建议优先使用 mmaphybrid 模式。

  • 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 参数说明如下。

参数

说明

spark.emr.serverless.excludedModules

esr-5.3.1 及以下需设置为 paimon,排除 EMR Serverless Spark 内置的 Paimon 模块,避免新旧类冲突。

spark.emr.serverless.user.defined.jars

自定义 JAR 的 OSS 地址。需包含与 Spark 版本匹配的 Paimon、paimon-eslib 及 ESLib/Lucene 依赖:可使用官方提供的 emr-3.5 适配 Jar,或用社区最新版本自行打包。

spark.sql.extensions

必须配置为 org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions。

table

目标 Paimon 表名。调用 oss.sys procedure 时填写 default.paimon_vector_demo。

index_column

写入同一个 es-index 的列。主向量列 emb_norm 必须放第一位,后续为过滤或返回的伴随标量列。

index_type

固定使用 es-index,生成当前 Elasticsearch paimon-store 可识别的索引文件。

options

传入分片行数及 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_namefile_sizebucketindex_typerow_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_pathoss_endpointoss_bucketoss_access_key_idoss_access_key_secret 为必填参数,缺失时接口会返回 Required fields 错误。

参数

是否必填

说明

auth_type

认证方式。示例使用 direct,即在请求中传入 OSS 访问凭据。生产环境建议优先使用更安全的托管凭据或 RAM 角色方式。

table_path

Paimon 表路径。

oss_endpoint

OSS endpoint。建议使用与 Elasticsearch 实例同地域的内网 endpoint。

oss_bucket

表所在 OSS bucket。

oss_access_key_id

访问 OSS 的 AccessKey ID。

oss_access_key_secret

访问 OSS 的 AccessKey Secret。

index_name

挂载后生成的 Elasticsearch 索引名称。

vector_field_name

用于向量检索的字段名称。

storage_mode

索引读取模式,例如 mmap、hybrid 或 heap。可用值以当前版本为准。

source_enabled

是否启用 Paimon source 回源。

return_fields

需要从 Paimon 表中返回的字段。

timeout

mount 超时时间,单位为秒。

建议首次验证使用 mmaphybrid。如果需要使用 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=768similarity=dot_product

  • 标量字段 idlabelcontent_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
  }
}
说明:挂载索引的标量列(idlabelcontent_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=truedata-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 数和副本数评估内存。