You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

本地K8s部署的Kubeflow容器间数据框与模型传递方法咨询

Kubeflow v2容器间传递DataFrame与模型的实现方案

一、使用func_to_container_op快速实现(无需自定义Dockerfile)

1. 传递DataFrame

用Kubeflow v2的Input/Output配合Dataset工件序列化DataFrame,推荐用Parquet格式(比CSV更高效):

from kfp.v2 import dsl
from kfp.v2.dsl import Input, Output, Dataset
import pandas as pd

@dsl.func_to_container_op
def load_data(output_df: Output[Dataset]):
    # 读取本地数据(本地文件需在镜像内或通过Volume挂载,下文讲处理方式)
    df = pd.read_csv("/path/to/local/data.csv")
    # 将DataFrame序列化到Output指定路径
    df.to_parquet(output_df.path)

@dsl.func_to_container_op
def process_data(input_df: Input[Dataset], output_processed: Output[Dataset]):
    # 读取上游传递的DataFrame
    df = pd.read_parquet(input_df.path)
    # 数据处理示例
    df["new_feature"] = df["old_feature"] * 1.5
    # 输出处理后的DataFrame
    df.to_parquet(output_processed.path)

@dsl.pipeline(name="dataframe-transfer-pipeline")
def pipeline():
    load_task = load_data()
    process_task = process_data(input_df=load_task.outputs["output_df"])

2. 传递模型

用Model工件保存训练好的模型,支持Pickle、HDF5、ONNX等格式:

from kfp.v2.dsl import Model
import joblib
from sklearn.ensemble import RandomForestClassifier

@dsl.func_to_container_op
def train_model(input_df: Input[Dataset], output_model: Output[Model]):
    df = pd.read_parquet(input_df.path)
    X = df.drop("label", axis=1)
    y = df["label"]
    # 训练模型示例
    model = RandomForestClassifier(n_estimators=100)
    model.fit(X, y)
    # 保存模型到Output路径
    joblib.dump(model, output_model.path)

@dsl.func_to_container_op
def predict(input_df: Input[Dataset], input_model: Input[Model]):
    df = pd.read_parquet(input_df.path)
    model = joblib.load(input_model.path)
    predictions = model.predict(df)
    # 处理预测结果逻辑
    print(predictions[:10])

本地文件的处理方式

如果要让容器访问本地文件,有两种常用方案:

  • 挂载PersistentVolume (PV):在Pipeline中定义Volume,将本地目录挂载到容器指定路径:
from kfp.v2.dsl import VolumeOp

@dsl.pipeline(name="local-file-pipeline")
def pipeline():
    # 创建PV(需集群有可用StorageClass)
    volume_op = VolumeOp(
        name="create-local-volume",
        size="1Gi",
        storage_class="standard"
    )
    # 将PV挂载到load_data任务的容器目录
    load_task = load_data().add_pvolumes({"/path/to/local": volume_op.outputs["volume"]})
  • 打包到镜像:如果文件是固定不变的,可通过func_to_container_op的base_image参数,指定包含该文件的自定义镜像(镜像构建时用COPY指令把本地文件复制进去)。

二、自定义Dockerfile构建容器实现

如果需要定制运行环境,可自行构建Docker镜像并定义Kubeflow组件:

1. 编写Dockerfile

FROM python:3.9-slim

# 安装依赖
RUN pip install pandas scikit-learn joblib kfp

# 复制本地数据文件到镜像(可选,适用于固定文件场景)
COPY ./local_data.csv /app/data.csv

# 复制组件代码
COPY ./load_data.py /app/load_data.py

WORKDIR /app

2. 编写组件代码(load_data.py)

import argparse
import pandas as pd

def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--output_df", type=str, required=True)
    args = parser.parse_args()
    
    df = pd.read_csv("/app/data.csv")
    df.to_parquet(args.output_df)

if __name__ == "__main__":
    main()

3. 定义并使用组件

from kfp.v2 import dsl
from kfp.v2.components import load_component_from_file

# 加载自定义组件(需提前编写component.yaml)
load_data_component = load_component_from_file("component.yaml")
# component.yaml示例:
# name: Load Local Data
# inputs: []
# outputs:
#   - name: output_df
#     type: Dataset
# implementation:
#   container:
#     image: your-registry/your-custom-image:v1
#     command: ["python", "/app/load_data.py"]
#     args: ["--output_df", {outputPath: output_df}]

@dsl.pipeline(name="custom-docker-pipeline")
def pipeline():
    load_task = load_data_component()
    # 后续任务通过Input/Output接收传递的工件,逻辑同func_to_container_op方式

本地文件传递(自定义镜像场景)

  • PV挂载:和func_to_container_op方式一致,在任务中用add_pvolumes挂载本地目录到容器路径。
  • 镜像打包:直接在Dockerfile中用COPY指令把本地文件复制到镜像内,容器内直接访问对应路径即可。

内容的提问来源于stack exchange,提问作者minio999

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 18:31:49