本文介绍如何基于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相似度检索。
方案整体链路包含图库摄入、查询图片摄入和向量检索三个部分。将图库与查询图片分别写入两张Paimon表,可避免向量检索时召回查询图片自身,并便于分别管理索引、生命周期与访问权限。
本文以TensorFlow公开的Flower Photos数据集(约3670张JPG图片,含daisy、dandelion、roses、sunflowers、tulips 5类)作为图库图片,另外准备5张数据集之外的花卉图片作为查询图片。运行后每张查询图片从图库中召回相似度Top-3的图片,共15条结果。
前提条件
-
已开通实时计算Flink版并创建工作空间,实时计算引擎版本为VVR 11.8.0及以上,详情请参见开通实时计算Flink版。
-
已创建OSS Bucket,且与Flink工作空间处于同一地域,详情请参见创建存储空间。
-
已开通轻量消息队列(原MNS)并创建队列,且与OSS Bucket处于同一地域,详情请参见开通轻量消息队列(原 MNS)。。
-
已开通DLF并创建Paimon类型的Catalog,详情请参见快速使用DLF。
-
已了解Flink AI服务的使用方式,详情请参见Flink AI服务(内置模型)。
-
已下载测试数据。图库图片使用Flower Photos数据集(约218MB),查询图片使用数据集之外的花卉照片,每类花卉各准备1张,共5张。
使用限制
-
仅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路径 |
说明 |
|
图库图片 |
|
存放Flower Photos全部约3670张图片 |
|
查询图片 |
|
存放5张数据集之外的花卉查询图片 |
将Flower Photos解压后的5个子目录(daisy、dandelion、roses、sunflowers、tulips)上传至图库目录,将准备的5张查询图片上传至查询目录。
2. 创建MNS队列
-
在左侧导航栏,选择队列模型 > 队列列表,选择目标地域。
-
单击创建队列,配置队列参数后单击确定。
|
参数 |
说明 |
|
名称 |
队列名称,例如 |
|
消息可见性超时时间 |
建议大于Flink Checkpoint间隔,避免消息重复投递。 |
|
消息保存时长 |
消息在队列中的最长存活时间,保持默认即可。 |
3. 配置OSS事件通知
为图库目录和查询目录分别创建事件通知规则,两条规则均投递到同一个MNS队列。
-
在MNS控制台左侧导航栏,单击事件通知。
-
在对象存储OSS页签,单击创建规则。
-
按下表配置图库目录规则,单击确定。
|
参数 |
说明 |
示例值 |
|
名称 |
规则名称。 |
|
|
事件类型 |
建议仅勾选创建、复制等对象创建事件。 |
ObjectCreated:\* |
|
匹配规则 |
前缀匹配。 |
|
|
接收终端 |
一对一订阅,输入步骤2中创建的队列名称。 |
|
重复以上操作创建查询目录规则,将匹配规则的前缀改为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表。
-
在实时计算控制台,创建数据摄入YAML作业草稿,详情请参见数据摄入YAML作业快速入门。
-
在编辑器中输入以下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连接器配置。 |
|
transform.projection |
使用 |
|
transform.table-options |
通过Paimon Blob列相关参数将图片二进制内容和URL纳入Blob管理, 本方案将图片写入Paimon Blob表存储。如果后续链路需要向外部服务(例如多模态模型服务)提供图片访问地址,可以使用Paimon提供的Blob预签名URL功能,详情请参见流式数据湖仓Paimon。 |
|
route |
将OSS CDC源表路由到Paimon目标表 |
|
sink |
Paimon Sink配置,通过DLF REST Catalog写入。 |
|
using.built-in-catalog |
可以直接复用已经创建的Catalog,详情请参见复用已有Catalog获取连接信息。 |
-
在部署详情的运行参数配置 > 其他配置中,追加步骤二中的OSS FileSystem鉴权参数。
-
保存并启动作业,作业名称建议设置为
oss-cdc-image-embedding-to-paimon。
作业启动后,OSS图库目录下约3670张图片会先通过全量扫描完成向量化入湖,之后新增图片将通过MNS事件通知实时进入image_search.image_assets表。
步骤四:创建查询图片摄入作业
复用步骤三的YAML模板,替换以下参数创建第二个作业:
|
参数 |
值 |
|
|
|
|
|
|
|
作业名称 |
|
其他配置项保持不变。作业启动后,5张查询图片的向量写入image_search.image_queries表。
步骤五:为图库表构建Global Index向量索引
在DLF控制台为image_search.image_assets表的embedding列构建Global Index向量索引,加速向量检索。
-
登录DLF控制台,进入目标Catalog。
-
找到
image_search.image_assets表,为embedding列创建向量类型的Global Index。 -
等待索引构建完成。索引构建完成前,
VECTOR_SEARCH函数无法在该向量列上执行。
详细操作请参见DLF构建Global Index。
步骤六:使用VECTOR_SEARCH完成以图搜图
两个摄入作业运行完成、Global Index构建就绪后,使用Flink SQL通过VECTOR_SEARCH函数完成相似度检索。
-
在实时计算控制台创建SQL流作业草稿。
-
在编辑器中输入以下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 |
运行参数。 |
VECTOR_SEARCH函数的完整语法请参见向量搜索。
步骤七:验证结果
作业启动后,5张查询图片分别从图库表召回Top-3相似图片,共输出15条结果。查询结果通过Print Connector输出到作业日志。
-
在实时计算控制台,进入目标作业的作业运维页面。
-
在作业日志页签,选择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等字段。