如何在Kubeflow各组件之间传递Python dataframe?
Kubeflow DataFrame跨组件传递方案
这个需求完全可以实现,你不需要手动编写持久化、路径传递的逻辑,Kubeflow(尤其是V2版本)原生支持结构化数据集类的输出传递,底层会自动完成序列化、存储、路径传递的全流程,也支持直接读取为DataFrame使用。
具体实现步骤
1. 定义输出类型为Dataset
在组件的函数参数中声明Output[Dataset]类型的输出,Kubeflow会自动为这个输出分配存储路径,你只需要把DataFrame写入这个路径对应的对象即可:
from kfp.v2.dsl import component, Output, Dataset import pandas as pd @component( base_image="python:3.9", packages_to_install=["pandas", "sqlalchemy"] # 按需添加依赖 ) def read_and_write(sql_query: str, output_df: Output[Dataset]): # 原有逻辑:执行查询生成df df = sql.to_dataframe() # 你的原查询转df逻辑 # 直接写入到输出Dataset的路径,Kubeflow自动处理存储 df.to_parquet(output_df.path, index=False) # 也可以用csv,但是parquet是二进制格式性能更好,保留类型信息 # df.to_csv(output_df.path, index=False)
这里优先推荐用parquet格式,不会出现序列化后的字段类型丢失问题,读写性能也远高于CSV。
2. 下游组件直接读取使用
下游组件只需要声明Input[Dataset]类型的输入参数,就可以直接读取为DataFrame,不需要手动处理路径:
from kfp.v2.dsl import component, Input, Dataset @component( base_image="python:3.9", packages_to_install=["pandas"] ) def process_df(input_df: Input[Dataset]): # 直接读取为DataFrame df = pd.read_parquet(input_df.path) # 后续处理逻辑 print(df.head())
3. 流水线中串联两个组件
组装流水线的时候直接把上游的输出参数传给下游的输入参数即可,不需要额外处理路径传递:
from kfp.v2 import dsl @dsl.pipeline(name="df-passing-demo") def pipeline(sql_query: str = "SELECT * FROM your_table"): read_task = read_and_write(sql_query=sql_query) process_task = process_df(input_df=read_task.outputs["output_df"])
补充说明
- 你不需要自己指定存储路径,Kubeflow会自动把输出存到你配置的流水线存储(一般是MinIO、GCS、S3这类对象存储)中,自动管理生命周期
- 如果是小体量的DataFrame(内存占用小于512KB),你也可以直接把DataFrame转成json字符串作为普通参数传递,但是大体积的数据集还是推荐用Dataset类型,避免触发参数传递的大小限制
- 除了parquet和csv,你也可以根据需求选择feather、pickle等序列化格式,只要上下游组件的读写逻辑对应即可
内容的提问来源于stack exchange,提问作者FrankNrg92
相关产品推荐
相关产品推荐

