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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 02:21:17