本文介绍如何在 EMR Serverless Spark 中加载 Paimon 自定义 JAR、配置 Spark SQL Session,并查询 Paimon Global Index。
前提条件
工作空间: 已创建 EMR Serverless Spark 工作空间。具体操作,请参见创建工作空间。
版本要求: 本文示例使用 Spark 3.5 和 Scala 2.12。
Catalog 配置: 已在工作空间中添加 DLF Paimon Catalog。具体操作,请参见Serverless Spark访问DLF。
索引准备: 目标表的 Global Index 已经生成并提交。具体操作,请参见构建 Global Index 实现查询加速和向量搜索。
权限配置: RAM 用户具备 DLF 相关 API 和数据权限。具体操作,请参见快速配置权限。
使用限制
样例限制:使用 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。
快速入门
进入 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 和多值索引。
查询
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 支持对应的谓词下推。使用 Bitmap 索引按时间段查询。
SELECT image_id, event_time, time_of_day, detected_objects FROM my_db.my_table WHERE time_of_day = 'night' LIMIT 20;使用多值索引查询包含指定对象的记录。
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按数组元素精确匹配。请使用常量作为待匹配元素,将字段与另一字段进行比较不等同于此处的常量谓词下推。组合查询夜间含车辆的记录。
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 会按照当前配置重新生成索引。