全部产品
Search
文档中心

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

更新时间:Sep 17, 2026

本文介绍如何在阿里云实时计算 Flink 版中注册内置 Paimon Catalog,并使用 VECTOR_SEARCH 函数查询已经创建向量索引的 Paimon 表。

前提条件

使用限制

  • 适用表类型: Global Index 仅适用于由 DLF 管理的 Paimon Append 表。向量字段的数据类型为 ARRAY<FLOAT>,同一字段中的向量维度必须一致。

  • 上游数据流: VECTOR_SEARCH 仅支持处理非更新流,即只包含 INSERT 类型消息的数据流。Append 表的读取会产生非更新流;主键表的 Changelog 流和 CDC 源包含 UPDATE、DELETE 消息,不能作为 COLUMN_TO_QUERY 的来源。

  • 索引覆盖范围: 默认只有已被向量索引覆盖的行参与检索。索引提交之后新写入的数据,在下一次索引提交前无法被检索,查询也不会报错。

准备工作

说明

本文后续操作均基于 DLF 共享样例数据集。示例在 Flink 中将该数据集对应的 DLF 只读 Catalog 注册为 dlf_samples,并查询其中的样例表。如果使用业务数据,请将示例中的 Catalog、数据库、表和字段名替换为实际值。

在 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>]
)

输入参数

参数

数据类型

说明

SEARCH_TABLE

TABLE

被检索的 Paimon 表,以 TABLE 关键字引导三段式表名。可以与上游数据流为同一张表,也可以是不同表。

COLUMN_TO_SEARCH

DESCRIPTOR

被检索表上已提交向量索引的向量列,以 DESCRIPTOR 关键字引导列名。同一向量列上只能存在一种向量索引类型。

COLUMN_TO_QUERY

ARRAY<FLOAT>/ARRAY<DOUBLE>

上游记录中用于发起检索的向量列。

TOP_K

INT

每条上游记录返回的相似结果条数上限。

CONFIG

MAP<STRING,STRING>

可配置的运行参数,分为引擎运行参数和 Paimon 向量索引查询参数两类。

VECTOR_SEARCH 的完整语法和输入参数说明,请参见向量搜索(VECTOR_SEARCH)。

CONFIG 参数

CONFIG 中的配置对上游每一条记录的检索都生效。

引擎运行参数控制检索的执行方式:

参数

数据类型

默认值

说明

async

Boolean

无

是否启用异步模式。默认情况下,引擎根据连接器支持的模式选择执行模式;连接器同时支持异步与同步模式时,优先使用异步模式以提升吞吐。如果连接器不支持指定的模式,作业会报错。

max-concurrent-operations

Integer

10

异步模式下的最大并行请求数。

output-mode

Enum

ORDERED

异步操作的输出模式,可选 ORDERED、ALLOW_UNORDERED。

timeout

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 按照当前配置自动重新生成索引。