Vertex AI Pipelines条件跳过步骤时依赖输出报错求助
解决Vertex AI Pipelines中条件步骤输出缺失的问题
在使用KFP(v1.8.11)和Vertex AI Pipelines时,当通过dsl.Condition跳过task1步骤后,无条件运行的task2因依赖task1的输出而编译失败,报错提示找不到task1-output_path。这是因为编译器无法识别未运行步骤的输出,需要为这种场景提供默认的输出占位。
解决方案
核心思路是新增一个生成默认输入的组件,在do_task1为False时运行该组件替代task1,为task2提供合法输入;通过条件分支选择使用task1或默认组件的输出传递给task2。
修改后的完整代码
from typing import NamedTuple from kfp import dsl from kfp.v2.dsl import Dataset, Input, OutputPath, component from kfp.v2 import compiler from google.cloud.aiplatform import pipeline_jobs @component( base_image="python:3.9", packages_to_install=["pandas"] ) def task1( task1_name: str, output_path: OutputPath("Dataset"), ) -> NamedTuple("Outputs", [("output_1", str), ("output_2", int)]): import pandas as pd output_1 = task1_name + "-processed" output_2 = 2 df_output_1 = pd.DataFrame({"output_1": [output_1]}) df_output_1.to_csv(output_path, index=False) return (output_1, output_2) @component( base_image="python:3.9", packages_to_install=["pandas"] ) def task2( task1_output: Input[Dataset], ) -> str: import pandas as pd task1_input = pd.read_csv(task1_output.path).values[0][0] return task1_input # 新增默认组件:生成task2所需的默认Dataset和输出值 @component( base_image="python:3.9", packages_to_install=["pandas"] ) def get_default_task1_output( output_path: OutputPath("Dataset"), ) -> NamedTuple("Outputs", [("output_1", str), ("output_2", int)]): import pandas as pd # 根据业务需求调整默认值 default_output_1 = "default-processed" default_output_2 = 0 df_default = pd.DataFrame({"output_1": [default_output_1]}) df_default.to_csv(output_path, index=False) return (default_output_1, default_output_2) @dsl.pipeline( pipeline_root='pipeline_root', name='pipelinename', ) def pipeline( do_task1: bool, task1_name: str, ): # 分条件运行task1或默认组件 with dsl.Condition(do_task1 == True, name="run_task1"): task1_op = task1(task1_name=task1_name) with dsl.Condition(do_task1 == False, name="use_default"): default_op = get_default_task1_output() # 选择符合当前条件的输出传递给task2 task2_input = dsl.one_of(task1_op.outputs["output_path"], default_op.outputs["output_path"]) task2_op = task2(task1_output=task2_input) if __name__ == '__main__': do_task1 = False # 可修改为True/False进行测试 # 编译流水线 compiler.Compiler().compile( pipeline_func=pipeline, package_path='pipeline.json') # 创建流水线运行任务 pipeline_run = pipeline_jobs.PipelineJob( display_name='pipeline-display-name', pipeline_root='pipelineroot', job_id='pipeline-job-id', template_path='pipeline.json', parameter_values={ 'do_task1': do_task1, 'task1_name': 'Task 1', }, enable_caching=False ) # 执行流水线 pipeline_run.run()
关键修改说明
- 新增
get_default_task1_output组件:生成task2所需的默认Dataset和输出参数,确保条件不满足时仍有合法输入。 - 双条件分支:分别处理
do_task1为True和False的场景,保证无论哪种情况都有对应的输出组件实例化。 - 使用
dsl.one_of:让编译器识别分支选择逻辑,自动匹配当前条件下的有效输出,避免因输出缺失导致的编译错误。
内容的提问来源于stack exchange,提问作者VanDough
相关产品推荐
相关产品推荐

