全部产品
Search
文档中心

实时计算Flink版:AI全模态入湖实时以图搜图

更新时间:Aug 14, 2026

本文介绍如何基于Flink OSS CDC连接器、Flink AI服务与DLF Paimon构建实时以图搜图方案,实现图片实时发现、向量化入湖与Flink SQL相似度检索的完整链路。

适用场景

本方案适用于以OSS为图片存储、需要实时或准实时完成图片相似度检索的业务场景,典型包括:

  • 电商拍照找相似商品:用户上传商品图片,从持续更新的商品图库中召回外观、颜色、形状和风格相近的商品,无需依赖文件名或人工标签。生产系统可结合类目、品牌、价格、库存等属性对向量召回结果进行二次过滤与排序。

  • 企业素材库与专业图库检索:设计师上传参考图,从OSS素材库中查找风格、主体、颜色或构图相近的历史素材。同样适用于植物、动物、文物、工业零部件等专业图库。

  • 商品重复上架治理:使用已上架商品图片检索图库,发现可能重复发布或高度相似的商品。图片Embedding偏向语义与视觉相似度,若需识别仅经过缩放或压缩的完全相同文件,可结合MD5、感知哈希等技术作为补充。

方案

图片检索场景中,图片通常存放在对象存储OSS,业务上要求新增图片能够被及时感知、向量化并进入检索链路。基于传统人工标签或批量导入方式,存在感知延迟高、标注成本高、多套存储数据一致性难以保障等问题。

本方案通过以下组件协同实现图片实时入湖与向量检索:

  • OSS:承载原始图片文件,通过事件通知向MNS队列投递对象创建事件。

  • MNS:接收OSS事件通知,供Flink OSS CDC连接器实时消费。

  • Flink OSS CDC连接器:实时消费MNS事件消息,感知OSS新增图片。

  • Flink AI服务:通过AI_IMAGE_EMBED函数调用内置qwen3-vl-embedding多模态模型,将图片转换为2560维向量。

  • DLF Paimon:统一存储图片元数据与向量,并通过Global Index为向量列构建索引。

  • VECTOR_SEARCH函数:在Flink SQL中基于Paimon表完成Top-K相似度检索。

image

方案整体链路包含图库摄入、查询图片摄入和向量检索三个部分。将图库与查询图片分别写入两张Paimon表,可避免向量检索时召回查询图片自身,并便于分别管理索引、生命周期与访问权限。

本文以TensorFlow公开的Flower Photos数据集(约3670张JPG图片,含daisy、dandelion、roses、sunflowers、tulips 5类)作为图库图片,另外准备5张数据集之外的花卉图片作为查询图片。运行后每张查询图片从图库中召回相似度Top-3的图片,共15条结果。

前提条件

使用限制

  • 仅VVR 11.8.0及以上版本支持该方案。

  • AI_IMAGE_EMBED函数的吞吐量受Flink AI服务限流限制,超出配额时作业可能出现反压。每个主账号每个地域每个自然月前100万tokens免费,超出部分按量付费,详情请参见AI服务计费。

  • VECTOR_SEARCH函数依赖Paimon表上的DLF Global Index,索引构建完成后新增数据需等待索引更新才能被检索。

  • 为保证检索结果稳定,同一张图片使用唯一且不重复的Object Key;图片更新时上传为新对象。

步骤一:准备OSS Bucket与MNS队列

1. 规划OSS图片目录

在目标OSS Bucket中规划两个虚拟目录,分别存放图库图片和查询图片。

用途

OSS路径

说明

图库图片

oss://<yourBucketName>/cdc-image-search/images/flowers/

存放Flower Photos全部约3670张图片

查询图片

oss://<yourBucketName>/cdc-image-search/queries/flowers/

存放5张数据集之外的花卉查询图片

将Flower Photos解压后的5个子目录(daisy、dandelion、roses、sunflowers、tulips)上传至图库目录,将准备的5张查询图片上传至查询目录。

2. 创建MNS队列

  1. 登录轻量消息队列(原 MNS)控制台。

  2. 在左侧导航栏,选择队列模型 > 队列列表,选择目标地域。

  3. 单击创建队列,配置队列参数后单击确定。

参数

说明

名称

队列名称,例如cdc-image-search-queue。

消息可见性超时时间

建议大于Flink Checkpoint间隔,避免消息重复投递。

消息保存时长

消息在队列中的最长存活时间,保持默认即可。

3. 配置OSS事件通知

为图库目录和查询目录分别创建事件通知规则,两条规则均投递到同一个MNS队列。

  1. 在MNS控制台左侧导航栏,单击事件通知。

  2. 在对象存储OSS页签,单击创建规则。

  3. 按下表配置图库目录规则,单击确定。

参数

说明

示例值

名称

规则名称。

image-assets-rule

事件类型

建议仅勾选创建、复制等对象创建事件。

ObjectCreated:\*

匹配规则

前缀匹配。

cdc-image-search/images/flowers/

接收终端

一对一订阅,输入步骤2中创建的队列名称。

cdc-image-search-queue

重复以上操作创建查询目录规则,将匹配规则的前缀改为cdc-image-search/queries/flowers/,规则名称改为image-queries-rule。

4. 记录连接信息

在队列列表页面,单击目标队列名称进入详情页,记录以下信息供后续Flink作业使用:

  • Endpoint:与Flink处于同一地域时使用VPC地址,格式为http://<account-id>.mns.<region>-internal.aliyuncs.com。

  • 队列名称:创建队列时设置的名称。

步骤二:配置作业访问OSS的鉴权参数

FETCH_CONTENT函数通过Flink FileSystem读取OSS图片内容。除OSS CDC连接器自身的凭据外,还需在作业运行参数中配置OSS文件系统访问信息。

登录实时计算控制台,在运维中心 > 作业运维 > 部署详情 > 运行参数配置 > 其他配置中添加以下参数:

fs.oss.bucket.<yourBucketName>.accessKeyId: <yourAccessKeyId>
fs.oss.bucket.<yourBucketName>.accessKeySecret: <yourAccessKeySecret>

将<yourBucketName>替换为目标OSS Bucket名称,AccessKey须具备ListObjects和GetObject权限。详情请参见配置Bucket鉴权信息。

运行身份还需具备以下权限:

  • MNS队列消息消费和删除权限。

  • DLF Catalog和Paimon Warehouse读写权限。

  • Flink AI服务模型调用权限。

步骤三:创建图库图片摄入作业

创建数据摄入YAML作业草稿,将OSS图库图片实时向量化并写入Paimon表。

  1. 在实时计算控制台,创建数据摄入YAML作业草稿,详情请参见数据摄入YAML作业快速入门。

  2. 在编辑器中输入以下YAML内容。

source:
  type: oss-cdc
  name: OSS CDC Source
  endpoint: http://<yourAccountId>.mns.<yourRegion>-internal.aliyuncs.com
  region: <yourRegion>
  queue-name: cdc-image-search-queue
  access-key-id: ${secret_values.mns_ak_id}
  access-key-secret: ${secret_values.mns_ak_secret}
  scan.startup.mode: INITIAL
  oss-endpoint: https://oss-<yourRegion>-internal.aliyuncs.com
  oss-bucket: <yourBucketName>
  path: cdc-image-search/images/flowers/

transform:
  - source-table: <yourBucketName>
    projection: >-
      key, url, fileName, eTag,
      FETCH_CONTENT(`url`) AS blob_field,
      AI_IMAGE_EMBED(FETCH_CONTENT(`url`), 'qwen3-vl-embedding') AS embedding
    table-options: blob-field=url,blob_field;blob-descriptor-field=url;blob-as-descriptor=true;row-tracking.enabled=true;data-evolution.enabled=true

route:
  - source-table: <yourBucketName>
    sink-table: image_search.image_assets

sink:
  type: paimon
  name: Paimon Sink
  using.built-in-catalog: <yourCatalog>

配置段

说明

source

OSS CDC连接器配置。scan.startup.mode设为INITIAL表示先全量扫描OSS指定路径下的存量图片,再切换为增量消费MNS事件。path用于限定扫描前缀。详情请参见OSS CDC。

transform.projection

使用FETCH_CONTENT函数读取OSS图片二进制内容,使用AI_IMAGE_EMBED函数调用qwen3-vl-embedding模型生成2560维向量。详情请参见Flink AI服务(内置模型)。

transform.table-options

通过Paimon Blob列相关参数将图片二进制内容和URL纳入Blob管理,row-tracking.enabled和data-evolution.enabled用于支撑后续Global Index构建。

本方案将图片写入Paimon Blob表存储。如果后续链路需要向外部服务(例如多模态模型服务)提供图片访问地址,可以使用Paimon提供的Blob预签名URL功能,详情请参见流式数据湖仓Paimon。

route

将OSS CDC源表路由到Paimon目标表image_search.image_assets。

sink

Paimon Sink配置,通过DLF REST Catalog写入。

using.built-in-catalog

可以直接复用已经创建的Catalog,详情请参见复用已有Catalog获取连接信息。

  1. 在部署详情的运行参数配置 > 其他配置中,追加步骤二中的OSS FileSystem鉴权参数。

  2. 保存并启动作业,作业名称建议设置为oss-cdc-image-embedding-to-paimon。

作业启动后,OSS图库目录下约3670张图片会先通过全量扫描完成向量化入湖,之后新增图片将通过MNS事件通知实时进入image_search.image_assets表。

步骤四:创建查询图片摄入作业

复用步骤三的YAML模板,替换以下参数创建第二个作业:

参数

值

source.path

cdc-image-search/queries/flowers/

route.sink-table

image_search.image_queries

作业名称

oss-cdc-query-image-embedding-to-paimon

其他配置项保持不变。作业启动后,5张查询图片的向量写入image_search.image_queries表。

步骤五:为图库表构建Global Index向量索引

在DLF控制台为image_search.image_assets表的embedding列构建Global Index向量索引,加速向量检索。

  1. 登录DLF控制台,进入目标Catalog。

  2. 找到image_search.image_assets表,为embedding列创建向量类型的Global Index。

  3. 等待索引构建完成。索引构建完成前,VECTOR_SEARCH函数无法在该向量列上执行。

详细操作请参见DLF构建Global Index。

步骤六:使用VECTOR_SEARCH完成以图搜图

两个摄入作业运行完成、Global Index构建就绪后,使用Flink SQL通过VECTOR_SEARCH函数完成相似度检索。

  1. 在实时计算控制台创建SQL流作业草稿。

  2. 在编辑器中输入以下SQL。

CREATE TEMPORARY TABLE print_sink (
  query_key   VARCHAR,
  matched_key VARCHAR,
  score       DOUBLE
) WITH (
  'connector' = 'print'
);

INSERT INTO print_sink
SELECT
  q.`key`  AS query_key,
  vs.`key` AS matched_key,
  vs.score
FROM `<yourCatalog>`.image_search.image_queries AS q,
LATERAL TABLE (
  VECTOR_SEARCH(
    SEARCH_TABLE    => TABLE `<yourCatalog>`.image_search.image_assets,
    COLUMN_TO_SEARCH => DESCRIPTOR(embedding),
    COLUMN_TO_QUERY  => q.embedding,
    TOP_K            => 3,
    CONFIG           => MAP['async', 'false']
  )
) AS vs;

将<yourCatalog>替换为实际的DLF Catalog名称。关键参数说明如下:

参数

说明

SEARCH_TABLE

图库表,包含全部图库图片的向量。

COLUMN_TO_SEARCH

图库表中已建立Global Index的向量列。

COLUMN_TO_QUERY

查询表中每一行的查询向量。

TOP_K

每张查询图片返回的相似图片数量。

CONFIG

运行参数。async为false时使用同步查询模式。

VECTOR_SEARCH函数的完整语法请参见向量搜索。

步骤七:验证结果

作业启动后,5张查询图片分别从图库表召回Top-3相似图片,共输出15条结果。查询结果通过Print Connector输出到作业日志。

  1. 在实时计算控制台,进入目标作业的作业运维页面。

  2. 在作业日志页签,选择Job Manager > Stdout,搜索PrintSinkOutputWriter相关日志。

输出示例:

+I[queries/flowers/daisy-01.jpg, images/flowers/daisy/100080576_f52e8ee070_n.jpg, 0.912]
+I[queries/flowers/daisy-01.jpg, images/flowers/daisy/10140303196_b88d3d6cec.jpg, 0.895]
+I[queries/flowers/daisy-01.jpg, images/flowers/daisy/10172379554_b296050f82_n.jpg, 0.881]

查询路径和匹配路径分属两张不同的Paimon表,检索时不会召回查询图片自身。

若需在结果中直接查看图片URL或文件名,可在print_sink中扩展oss_url、file_name等字段。