全部产品
Search
文档中心

数据湖构建:使用 Spark 查询 Paimon Global Index

更新时间:Sep 17, 2026

本文介绍如何在 EMR Serverless Spark 中加载 Paimon 自定义 JAR、配置 Spark SQL Session,并查询 Paimon Global Index。

前提条件

使用限制

  • 样例限制:使用 DLF 共享样例表时,仅可体验已创建的索引。共享样例 Catalog 为只读,不支持在样例表上创建业务索引。

  • 运行依赖:执行 Bitmap 和多值索引查询前,需确认所用 Paimon JAR 支持 Bitmap、多值索引及数组谓词下推。

  • 索引验证:查询成功不代表索引已经生效,需结合执行计划和扫描指标确认索引效果。

准备 Spark 环境

步骤一:下载 Paimon JAR

下载 paimon-ali-emr-spark-3.5-2-ali-2.5-vector-all-in-one.jar。该 JAR 包含 Paimon Spark Connector、DLF Catalog、Jindo 文件系统和向量索引运行依赖。

步骤二:上传 JAR

将 JAR 上传到 EMR Serverless Spark 文件管理,或者上传到与工作空间同地域的 OSS Bucket。上传完成后,记录文件的 OSS 地址。

具体操作,请参见管理文件。

步骤三:配置 Spark SQL Session

本文使用 Spark SQL Session 运行后续示例。在 EMR Serverless Spark 工作空间中创建或选择 Spark SQL Session(详见会话管理),并增加以下配置:

spark.emr.serverless.excludedModules    paimon
spark.emr.serverless.user.defined.jars  oss://<bucket>/jars/paimon-ali-emr-spark-3.5-2-ali-2.5-vector-all-in-one.jar
spark.sql.extensions                    org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions

其中,spark.emr.serverless.excludedModules=paimon 用于排除运行环境内置的 Paimon 版本,避免与自定义 JAR 冲突。

修改 JAR 或 Session 参数后,需要重新启动 Session。

快速入门

说明

可以使用共享样例数据集创建只读 Catalog,体验 B-Tree 查询、IVF-PQ 向量检索和混合检索。本节使用 DLF 共享样例表 search_samples.berkeley_deepdrive_100k_images 进行演示。

进入 EMR Serverless Spark 工作空间的数据开发页面,在 SparkSQL 编辑器中选择已配置的 Spark SQL Session、实际使用的 DLF Paimon Catalog 和 search_samples 数据库,然后运行以下 SQL。具体操作,请参见SparkSQL开发。

查看 Global Index

查询 Paimon table_indexes 系统表,确认 B-Tree 和 IVF-PQ 索引已经提交:

SELECT
    index_type,
    index_field_name,
    COUNT(*) AS index_file_count,
    SUM(row_count) AS indexed_rows
FROM `berkeley_deepdrive_100k_images$table_indexes`
WHERE index_type IN ('btree', 'ivf-pq')
GROUP BY index_type, index_field_name
ORDER BY index_type, index_field_name;

查询结果中应包含以下信息:

  • weather 对应的 btree 索引。

  • image_embedding 对应的 ivf-pq 索引。

  • indexed_rows 大于 0,表示索引中已经包含可查询的数据。

使用 B-Tree 按天气查询

B-Tree 索引会自动用于 weather 等值条件,无需在 SQL 中指定索引名称。

SELECT
    image_id,
    event_time,
    weather,
    time_of_day,
    road_scene,
    detected_objects
FROM berkeley_deepdrive_100k_images
WHERE weather = 'rainy'
LIMIT 20;

样例表的 B-Tree 索引建在 weather 字段上,因此本示例使用 weather 等值条件。

使用 IVF-PQ 查询相似图片

以下 SQL 使用样例表中一张图片的 image_embedding 查询相似图片 Top 10,无需在 SQL 中填写完整的 768 维向量。

WITH query_image AS (
    SELECT image_id, image_embedding
    FROM berkeley_deepdrive_100k_images
    WHERE image_id = 'c6459d6c-2a477f42'
)
SELECT
    q.image_id AS query_image_id,
    r.image_id,
    r.event_time,
    r.weather,
    r.time_of_day,
    r.road_scene,
    r.detected_objects,
    r.__paimon_search_score
FROM query_image q,
LATERAL (
    SELECT
        image_id,
        event_time,
        weather,
        time_of_day,
        road_scene,
        detected_objects,
        __paimon_search_score
    FROM vector_search(
        'berkeley_deepdrive_100k_images',
        'image_embedding',
        q.image_embedding,
        10
    )
) r
ORDER BY r.__paimon_search_score DESC;

__paimon_search_score 表示相似度得分,值越大表示越相似。

使用 B-Tree 和 IVF-PQ 进行混合检索

混合检索先通过结构化条件限定数据范围,再在该范围内执行向量召回。结构化条件需要写在 vector_search 子查询内。

WITH query_image AS (
    SELECT image_id, image_embedding
    FROM berkeley_deepdrive_100k_images
    WHERE image_id = 'c6459d6c-2a477f42'
)
SELECT
    q.image_id AS query_image_id,
    r.image_id,
    r.event_time,
    r.weather,
    r.time_of_day,
    r.road_scene,
    r.detected_objects,
    r.__paimon_search_score
FROM query_image q,
LATERAL (
    SELECT
        image_id,
        event_time,
        weather,
        time_of_day,
        road_scene,
        detected_objects,
        __paimon_search_score
    FROM vector_search(
        'berkeley_deepdrive_100k_images',
        'image_embedding',
        q.image_embedding,
        10
    )
    WHERE weather = 'rainy'
) r
ORDER BY r.__paimon_search_score DESC;

该查询会使用 weather 上的 B-Tree 索引缩小数据范围,再使用 image_embedding 上的 IVF-PQ 索引返回相似图片。

批量向量检索

以下 SQL 从共享样例表中读取 1000 张图片,并在同一张表中为每张图片检索 Top 3 相似结果。

WITH query_images AS (
    SELECT
        image_id,
        image_embedding
    FROM berkeley_deepdrive_100k_images
    WHERE image_embedding IS NOT NULL
    LIMIT 1000
)
SELECT /*+ REPARTITION(1) */
    q.image_id AS query_image_id,
    r.image_id AS result_image_id,
    r.event_time,
    r.weather,
    r.time_of_day,
    r.road_scene,
    r.detected_objects,
    r.__paimon_search_score
FROM query_images q,
LATERAL (
    SELECT
        image_id,
        event_time,
        weather,
        time_of_day,
        road_scene,
        detected_objects,
        __paimon_search_score
    FROM vector_search(
        'berkeley_deepdrive_100k_images',
        'image_embedding',
        q.image_embedding,
        3
    )
) r
SORT BY query_image_id, r.__paimon_search_score DESC;

批量检索参数已有默认值,通常无需显式设置:

  • spark.paimon.vector-search.lateral-join.batch-size:默认值为 256。增大该值可减少索引加载和回表次数,但会增加内存占用与单批延迟。

  • spark.paimon.vector-search.lateral-join.parallelism:默认值为 16。提高该值可增加批量检索并发度,但需结合 Executor Core、存储 I/O 和索引缓存能力调整。

query_image_id 标识结果所属的查询图片,1000 张图片最多返回 3000 行。由于查询图片和目标图片来自同一张表,结果通常会包含图片自身。REPARTITION(1) 和 SORT BY 仅用于整理较小的测试结果,大规模输出时建议移除。

查询 Bitmap 和多值索引

本示例使用业务表 my_db.my_table。运行前,请确认该表已创建 Bitmap 和多值索引。

  1. 查询 table_indexes 系统表,确认索引已经提交。

    SELECT index_type, index_field_name, SUM(row_count) AS indexed_rows
    FROM `my_db`.`my_table$table_indexes`
    WHERE index_type IN ('bitmap', 'multivalue')
    GROUP BY index_type, index_field_name;

    查询结果应分别包含 bitmap/time_of_day 和 multivalue/detected_objects。索引文件已提交不代表查询一定命中索引,还需确认当前 Connector 支持对应的谓词下推。

  2. 使用 Bitmap 索引按时间段查询。

    SELECT image_id, event_time, time_of_day, detected_objects
    FROM my_db.my_table
    WHERE time_of_day = 'night'
    LIMIT 20;
  3. 使用多值索引查询包含指定对象的记录。

    SELECT image_id, event_time, time_of_day, detected_objects
    FROM my_db.my_table
    WHERE array_contains(detected_objects, 'car')
    LIMIT 20;

    array_contains 按数组元素精确匹配。请使用常量作为待匹配元素,将字段与另一字段进行比较不等同于此处的常量谓词下推。

  4. 组合查询夜间含车辆的记录。

    SELECT image_id, event_time, time_of_day, detected_objects
    FROM my_db.my_table
    WHERE time_of_day = 'night'
      AND array_contains(detected_objects, 'car')
    LIMIT 20;

可对上述 SELECT 语句执行 EXPLAIN EXTENDED,检查过滤条件是否下推,再结合扫描指标确认索引效果。仅查询成功不能证明索引生效。

传递向量查询参数

vector_search 的前四个参数依次为表名、向量字段、查询向量和 Top K。第五个参数用于传递可选的查询参数。

FROM vector_search(
    'berkeley_deepdrive_100k_images',
    'image_embedding',
    q.image_embedding,
    10,
    map('ivf.nprobe', '32')
)

未设置第五个参数时,Paimon 使用自动配置。通常可先使用自动配置,在召回率或查询耗时不满足要求时再调整参数。

带过滤条件的查询还可设置 ivf.max_initial_filter_expansion_factor。该参数仅在自动计算 ivf.nprobe 时生效,不能与显式设置的 ivf.nprobe 同时使用。批量 IVF-PQ 查询可设置 ivf_pq.batch_table_reuse。参数含义、默认值和调优方法,请参见配置 IVF-PQ 构建和查询参数。

其他 Spark 运行方式

本文示例通过 Spark SQL Session 执行。Spark 还支持其他作业开发和提交方式。无论采用哪种方式,都需要在运行环境中加载本文提供的 JAR,并配置能够访问目标 Paimon 表的 Catalog。具体操作,请参见开发管理。

附录:删除 Global Index

调优构建参数并需要重新生成索引时,可使用 Spark Procedure 删除已有索引。建议先使用 dry_run=true 查看命中的索引文件数量。

CALL paimon_catalog.sys.drop_global_index(
    table => 'search_samples.berkeley_deepdrive_100k_images',
    index_column => 'image_embedding',
    index_type => 'ivf-pq',
    dry_run => true
);

将 paimon_catalog 替换为实际 Catalog Name。确认目标表、字段和索引类型后,将 dry_run 设置为 false 再次执行。删除 Global Index 不会删除 Paimon 表数据。如果表中仍然开启对应索引的自动构建,DLF 会按照当前配置重新生成索引。