本文通过一个电商商品图片分析的示例,介绍如何使用 Flink Python DataFrame API 调用多模态大模型,完成从图片读取、预处理、生成商品描述到文本向量化的完整流程。DataFrame API 的更多介绍请参见功能概览。
前提条件
已创建 Flink 工作空间,详情请参见开通实时计算Flink版。
已开通 Flink AI 服务,详情请参见Flink AI服务(内置模型)。
(可选)本地开发环境配置:由于 Flink Python DataFrame API 为 VVR 11.8 及以上的专属功能,暂未在开源版 Flink 中提供,因此您需要使用下面的命令安装 VVR 版 PyFlink 依赖包:
pip install "ververica-flink==11.8.0"说明ververica-flink是API-Only包,仅提供类型定义和接口声明,用于本地开发,作业仍需在实时计算Flink版平台提交。
数据准备
在本例中,您需要准备测试图片文件(1.jpg、2.jpg、3.jpg),并通过文件管理功能上传到 Flink 工作空间。上传方法请参见文件管理。
作业开发
创建一个 Python 文件(如 pyflink_dataframe_example.py),按照以下步骤进行作业开发。您可在本文档末尾获取 Python 文件的完整代码。
步骤一:注册 AI 模型 Provider
使用 pf.set_model_provider注册两个 Model Provider:chat 用于图片理解,embedding 用于把商品描述转换为文本向量。
import pyflink.dataframe as pf
pf.set_model_provider(
"chat",
pf.OpenAICompatProvider(
task="chat/completions",
error_handling_strategy="RETRY",
retry_num=3,
retry_backoff_strategy="EXPONENTIAL",
),
)
pf.set_model_provider(
"embedding",
pf.OpenAICompatProvider(
task="embeddings",
error_handling_strategy="RETRY",
retry_num=3,
retry_backoff_strategy="EXPONENTIAL",
),
)步骤二:准备输入数据
使用 pf.from_dict 构造输入 DataFrame,包含商品 ID 和对应的 OSS 路径。
OSS_PREFIX = "<your_oss_prefix>"
df = pf.from_dict(
{
"image_id": ["1", "2", "3"],
"image_url": [
f"{OSS_PREFIX}/1.jpg",
f"{OSS_PREFIX}/2.jpg",
f"{OSS_PREFIX}/3.jpg",
],
}
)您可以根据 Flink 工作空间的存储类型确定 OSS_PREFIX 的值,详情请参见文件管理:
若存储类型为 OSS bucket,
OSS_PREFIX的值为oss://<您绑定的OSS Bucket名称>/artifacts/namespaces/<项目空间名称>;若存储类型为全托管存储,
OSS_PREFIX的值为oss://flink-fullymanaged-<工作空间ID>/artifacts/namespaces/<项目空间名称>。
在生产环境中,通常从外部数据源读取数据。DataFrame API 提供了多种数据读写函数,帮助您便捷地操作外部数据,请参见I/O 读写。
步骤三:使用多模态算子预处理图片
本例调用多模态算子,对图片 URL 执行以下操作:
通过
fetch_content()从 OSS 读取图片字节。使用
image.is_valid()校验图片。使用
image.decode()将编码字节解码。使用
image.size_filter()完成尺寸检查。使用
image.crop_black_border()、image.resize()和image.encode(output="data_url")生成适合多模态模型调用的 Data URL。
from pyflink.dataframe import col
preprocessed = (
df.with_column("image_bytes", col("image_url").fetch_content())
.with_column(
"image_valid",
col("image_bytes").image.is_valid(
pixel_limit=40_000_000,
concurrency=16,
),
)
.filter("image_valid")
.with_column(
"image",
col("image_bytes").image.decode(
on_error="null",
mode="RGB",
pixel_limit=40_000_000,
concurrency=16,
),
)
.drop_null(subset=["image"])
.with_columns(
min_size_pass=col("image").image.size_filter(
min_w=256,
min_h=256,
concurrency=16,
),
image_data=(
col("image")
.image.crop_black_border(concurrency=16)
.image.resize(width=1024, height=1024, concurrency=16)
.image.encode(
format="JPEG",
quality=85,
output="data_url",
concurrency=16,
)
),
)
.filter("min_size_pass")
)步骤四:调用 AI 函数
将预处理后的图片 Data URL 传给 llm.predict,使用 qwen3.6-plus 生成商品描述;再调用 llm.ai_embed,使用 text-embedding-v4 生成描述向量。
result = (
preprocessed.llm.predict(
"image_data",
provider="chat",
model="qwen3.6-plus",
output_type=pf.DataType.struct(
{"image_description": pf.DataType.string()}
),
content_type="IMAGE_URL",
system_prompt=(
"你是一个电商商品分析助手。请根据商品图片生成准确、简洁的商品描述,"
"只输出描述本身。"
),
temperature=0.2,
)
.llm.ai_embed(
"image_description",
provider="embedding",
model="text-embedding-v4",
)
.select(
"image_id",
"image_description",
"embedding",
)
)步骤五:输出结果
将处理结果写入 Parquet 文件。
result.write_parquet(
f"{OSS_PREFIX}/product_descriptions/"
)DataFrame API 的详细文档请参见API参考。
作业部署与运行
部署 Python 作业。部署流程请参见Flink Python作业。
配置作业:
在页面,单击目标作业名称。
在部署详情页签基础配置区域引擎版本,选择 vvr-11.8-jdk11-flink-1.20。
运行作业:
在页面,单击目标作业名称操作列中的启动。
选择无状态启动,单击启动,作业启动详情请参见作业启动。
单击启动后,作业状态变为运行中或已完成,则代表作业运行正常。如果您部署本文档Python测试文件,作业最终运行状态是已完成状态。
作业运行完成后,在文件管理页面下载
product_descriptions/目录下的 Parquet 文件,验证输出结果。预期输出包含以下列:列名
类型
说明
image_idSTRING
商品 ID
image_descriptionSTRING
模型生成的商品描述
embeddingARRAY<FLOAT>
描述文本对应的向量
完整代码
import pyflink.dataframe as pf
from pyflink.dataframe import col
OSS_PREFIX = "<your_oss_prefix>"
pf.set_model_provider(
"chat",
pf.OpenAICompatProvider(
task="chat/completions",
error_handling_strategy="RETRY",
retry_num=3,
retry_backoff_strategy="EXPONENTIAL",
),
)
pf.set_model_provider(
"embedding",
pf.OpenAICompatProvider(
task="embeddings",
error_handling_strategy="RETRY",
retry_num=3,
retry_backoff_strategy="EXPONENTIAL",
),
)
df = pf.from_dict(
{
"image_id": ["1", "2", "3"],
"image_url": [
f"{OSS_PREFIX}/1.jpg",
f"{OSS_PREFIX}/2.jpg",
f"{OSS_PREFIX}/3.jpg",
],
}
)
preprocessed = (
df.with_column("image_bytes", col("image_url").fetch_content())
.with_column(
"image_valid",
col("image_bytes").image.is_valid(
pixel_limit=40_000_000,
concurrency=16,
),
)
.filter("image_valid")
.with_column(
"image",
col("image_bytes").image.decode(
on_error="null",
mode="RGB",
pixel_limit=40_000_000,
concurrency=16,
),
)
.drop_null(subset=["image"])
.with_columns(
min_size_pass=col("image").image.size_filter(
min_w=256,
min_h=256,
concurrency=16,
),
image_data=(
col("image")
.image.crop_black_border(concurrency=16)
.image.resize(width=1024, height=1024, concurrency=16)
.image.encode(
format="JPEG",
quality=85,
output="data_url",
concurrency=16,
)
),
)
.filter("min_size_pass")
)
result = (
preprocessed.llm.predict(
"image_data",
provider="chat",
model="qwen3.6-plus",
output_type=pf.DataType.struct(
{"image_description": pf.DataType.string()}
),
content_type="IMAGE_URL",
system_prompt=(
"你是一个电商商品分析助手。请根据商品图片生成准确、简洁的商品描述,"
"只输出描述本身。"
),
temperature=0.2,
)
.llm.ai_embed(
"image_description",
provider="embedding",
model="text-embedding-v4",
)
.select(
"image_id",
"image_description",
"embedding",
)
)
result.write_parquet(
f"{OSS_PREFIX}/product_descriptions/"
)