本文介绍如何在阿里云实时计算 Flink 版中注册内置 Paimon Catalog,并使用 VECTOR_SEARCH 函数查询已经创建向量索引的 Paimon 表。
前提条件
权限配置: RAM 用户具备 DLF 相关 API 和数据权限。具体操作,请参见快速配置权限。
Flink 工作空间: 已开通阿里云实时计算 Flink 版工作空间,具体操作,请参见开通实时计算 Flink 版。VVR 版本为 11.9 或更高版本。
网络配置: Flink 工作空间 VPC 已加入 DLF 可信名单。具体操作,请参见可信 VPC 配置。
Catalog 配置: 已创建可访问目标 Paimon 表的 DLF Catalog,详情请参见数据目录。如果需要使用 DLF 共享样例数据集创建只读 Catalog,请参见使用 DLF 共享样例数据集。
向量索引: 目标表的向量列上已经生成并提交 Global Index。如果需要为业务表创建索引,请参见构建 Global Index 实现查询加速和向量搜索。
使用限制
适用表类型: Global Index 仅适用于由 DLF 管理的 Paimon Append 表。向量字段的数据类型为
ARRAY<FLOAT>,同一字段中的向量维度必须一致。上游数据流:
VECTOR_SEARCH仅支持处理非更新流,即只包含INSERT类型消息的数据流。Append 表的读取会产生非更新流;主键表的 Changelog 流和 CDC 源包含UPDATE、DELETE消息,不能作为COLUMN_TO_QUERY的来源。索引覆盖范围: 默认只有已被向量索引覆盖的行参与检索。索引提交之后新写入的数据,在下一次索引提交前无法被检索,查询也不会报错。
准备工作
在 Flink 中注册 Catalog
进入实时计算 Flink 版控制台,选择数据开发 > 数据查询,执行以下 SQL 创建 Catalog:
CREATE CATALOG `dlf_samples` WITH (
'type' = 'paimon',
'metastore' = 'rest',
'token.provider' = 'dlf',
'uri' = 'http://cn-hangzhou-vpc.dlf.aliyuncs.com',
'warehouse' = 'dlf_samples'
);参数说明:
uri:设置为http://<vpc_endpoint>,其中<vpc_endpoint>为 Flink 工作空间所在地域的 DLF VPC Endpoint。各地域的 Endpoint,请参见地域与服务接入点。warehouse:替换为目标表所在的实际 DLF Catalog Name。Catalog 名称:可以按需修改,示例中的 Catalog 名称为
dlf_samples。
创建完成后,刷新数据管理页面,确认 Catalog 中可以查看目标数据库和表。DLF Catalog 参数说明,请参见管理 Paimon Catalog。
查询向量索引
VECTOR_SEARCH 通过 LATERAL TABLE 与上游数据流关联。上游每流入一条记录,就以该记录的向量发起一次检索,命中的行持续追加到结果流。
语法
VECTOR_SEARCH(
SEARCH_TABLE => TABLE <search_table>,
COLUMN_TO_SEARCH => DESCRIPTOR(<column_to_search>),
COLUMN_TO_QUERY => <column_to_query>,
TOP_K => <top_k>[,
CONFIG => <config>]
)输入参数
参数 | 数据类型 | 说明 |
|
| 被检索的 Paimon 表,以 |
|
| 被检索表上已提交向量索引的向量列,以 |
|
| 上游记录中用于发起检索的向量列。 |
|
| 每条上游记录返回的相似结果条数上限。 |
|
| 可配置的运行参数,分为引擎运行参数和 Paimon 向量索引查询参数两类。 |
VECTOR_SEARCH 的完整语法和输入参数说明,请参见向量搜索(VECTOR_SEARCH)。
CONFIG 参数
CONFIG 中的配置对上游每一条记录的检索都生效。
引擎运行参数控制检索的执行方式:
参数 | 数据类型 | 默认值 | 说明 |
| Boolean | 无 | 是否启用异步模式。默认情况下,引擎根据连接器支持的模式选择执行模式;连接器同时支持异步与同步模式时,优先使用异步模式以提升吞吐。如果连接器不支持指定的模式,作业会报错。 |
| Integer | 10 | 异步模式下的最大并行请求数。 |
| Enum |
| 异步操作的输出模式,可选 |
| Duration | 3min | 从首次调用到异步操作最终完成的超时时间,可能包含多次重试,并在发生故障转移时被重置。 |
向量索引的检索语法与索引类型无关。Paimon 向量索引查询参数由 Paimon Connector 透传给向量索引后端,用于调整召回率与查询耗时。CONFIG 中的参数名由向量索引类型决定,以下示例中的 ivf.nprobe 适用于 IVF 系列索引:
VECTOR_SEARCH(
SEARCH_TABLE => TABLE <search_table>,
COLUMN_TO_SEARCH => DESCRIPTOR(<column_to_search>),
COLUMN_TO_QUERY => <column_to_query>,
TOP_K => 10,
CONFIG => MAP[
'async', 'true',
'ivf.nprobe', '32'
]
)通常不需要设置向量索引查询参数,Paimon 会使用默认配置。只有召回率或查询耗时不满足要求时,再调整相关参数。IVF-PQ 索引的参数含义和调优方法,请参见IVF-PQ 构建和查询参数。
返回结果
VECTOR_SEARCH 返回一张数据表。每条上游记录对应至多 TOP_K 行,每行包含被检索表的全部列,以及一个附加列 score。score 为 DOUBLE 类型,描述该行向量与输入向量的相似度,值越大表示越相似。
样例表使用内积作为距离度量。向量已归一化时,分数可以近似按余弦相似度理解。
示例
在实时计算 Flink 版的数据开发 > 数据查询页面中新建查询脚本。
以下示例使用 DLF 共享样例表 dlf_samples.search_samples.berkeley_deepdrive_100k_images 作为上游数据流,为每条记录检索 Top 10 相似图片。该表的 image_embedding 列上建有 IVF-PQ 向量索引。
SELECT
q.`image_id` AS query_image_id,
r.`image_id` AS similar_image_id,
r.`weather`,
r.`road_scene`,
r.`score`
FROM `dlf_samples`.`search_samples`.`berkeley_deepdrive_100k_images` AS q,
LATERAL TABLE(
VECTOR_SEARCH(
SEARCH_TABLE => TABLE `dlf_samples`.`search_samples`.`berkeley_deepdrive_100k_images`,
COLUMN_TO_SEARCH => DESCRIPTOR(image_embedding),
COLUMN_TO_QUERY => q.`image_embedding`,
TOP_K => 10
)
) AS r;实际业务中,查询向量通常来自消息队列或另一张持续写入的 Paimon 表。将上游替换为该数据流即可,VECTOR_SEARCH 部分不需要改动。
删除 Global Index
索引的构建与维护由 DLF 完成,删除已有索引需要通过 Flink Procedure 执行。调优构建参数并需要重新生成索引时,建议先使用 dry_run=true 查看命中的索引文件数量。
CALL dlf_samples.sys.drop_global_index(
`table` => 'search_samples.berkeley_deepdrive_100k_images',
`index_column` => 'image_embedding',
`index_type` => 'ivf-pq',
`dry_run` => true
);确认目标表、字段和索引类型后,将 dry_run 设置为 false 再次执行。删除操作不会删除 Paimon 表数据。如果表中仍然开启对应索引的自动构建,请勿手动调用 create_global_index,等待 DLF 按照当前配置自动重新生成索引。