如何在Kubeflow中获取流水线的输出值与生成的工件?
如何在Kubeflow中获取流水线的输出值与生成的工件?
我完全懂你的困扰——本地用DockerRunner时能直接通过res.output拿到结果,但在K8s集群上跑就没这么直观了。毕竟K8s上的流水线是异步执行的,结果存在Kubeflow的元数据存储里,得通过KFP Client的API主动拉取。下面分两部分给你具体的解决方法:
一、获取流水线的返回值
当你在K8s上启动流水线并等待完成后,run.wait_for_run_completion()返回的res只包含基础运行状态,要拿到具体的返回值,需要进一步调用API获取完整的运行详情。可以这么操作:
client = Client(host=host) # 注意:你代码里的kfp_pipline应该是pythagorean的笔误,这里修正过来 run = client.create_run_from_pipeline_func( pythagorean, arguments={"a": 3, "b": 4} ) res = run.wait_for_run_completion(timeout=3600) # 步骤1:获取完整的运行详情 run_details = client.get_run(run.run_id) # 步骤2:提取流水线的返回值 # 未指定输出名称时,默认键为"pipelineoutput";如果在@pipeline装饰器里指定了output_names,就用你定义的键 pipeline_return_value = run_details.run.outputs["pipelineoutput"].value print(f"流水线返回值: {pipeline_return_value}")
举个例子,如果你的流水线定义时指定了输出名称:
@dsl.pipeline(name='pythagorean', output_names=['hypotenuse']) def pythagorean(a: float, b: float) -> float: # 原有逻辑不变
那提取时就把键换成hypotenuse即可。
二、获取流水线生成的输出工件
如果你的流水线会生成输出工件(比如模型文件、数据集、日志文件等),可以通过两种方式获取它们的信息:
方法1:从运行详情中直接提取
# 从已获取的run_details中提取所有输出工件 output_artifacts = [] for output_key, output_value in run_details.run.outputs.items(): # 每个输出项可能关联多个工件 if output_value.artifacts: output_artifacts.extend(output_value.artifacts) # 遍历打印工件的关键信息 for artifact in output_artifacts: print(f"工件名称: {artifact.name}") print(f"集群内存储路径: {artifact.uri}") print(f"工件类型: {artifact.type.name}") print("---")
方法2:用Client的list_artifacts方法查询
# 直接查询该流水线运行产生的所有工件 artifact_list = client.list_artifacts(run_id=run.run_id) for artifact in artifact_list.artifacts: print(f"工件ID: {artifact.id}") print(f"存储URI: {artifact.uri}") print("---")
这里要注意:工件的URI一般是集群内部的存储路径(比如MinIO对象存储的路径),如果要在本地代码里下载处理,需要配置对应的存储客户端(比如MinIO的Python SDK)来访问。
为什么K8s和本地DockerRunner的获取方式不一样?
本地DockerRunner是在你的Python进程中同步执行流水线逻辑,返回的res是实际的PipelineTask对象,所以能直接拿到output。但在K8s上,流水线是提交到集群异步执行的,create_run_from_pipeline_func返回的run只是一个运行的引用,结果存在Kubeflow的元数据服务中,必须通过Client API主动拉取才能拿到。
备注:内容来源于stack exchange,提问作者Elad Project
相关产品推荐
相关产品推荐

