使用Flink Python DataFrame API进行电商商品图片分析

更新时间:
复制 MD 格式

本文通过一个电商商品图片分析的示例,介绍如何使用 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-flinkAPI-Only包,仅提供类型定义和接口声明,用于本地开发,作业仍需在实时计算Flink版平台提交。

数据准备

在本例中,您需要准备测试图片文件(1.jpg2.jpg3.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 执行以下操作:

  1. 通过 fetch_content() 从 OSS 读取图片字节。

  2. 使用 image.is_valid() 校验图片。

  3. 使用 image.decode() 将编码字节解码。

  4. 使用 image.size_filter() 完成尺寸检查。

  5. 使用 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参考

作业部署与运行

  1. 部署 Python 作业。部署流程请参见Flink Python作业

  2. 配置作业:

    1. 运维中心 > 作业运维页面,单击目标作业名称。

    2. 部署详情页签基础配置区域引擎版本,选择 vvr-11.8-jdk11-flink-1.20。

  3. 运行作业:

    1. 运维中心 > 作业运维页面,单击目标作业名称操作列中的启动

    2. 选择无状态启动,单击启动,作业启动详情请参见作业启动

      单击启动后,作业状态变为运行中已完成,则代表作业运行正常。如果您部署本文档Python测试文件,作业最终运行状态是已完成状态。

  4. 作业运行完成后,在文件管理页面下载 product_descriptions/ 目录下的 Parquet 文件,验证输出结果。预期输出包含以下列:

    列名

    类型

    说明

    image_id

    STRING

    商品 ID

    image_description

    STRING

    模型生成的商品描述

    embedding

    ARRAY<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/"
)