如何实现第一个Vertex AI流水线成功完成后自动运行第二个?
实现Vertex AI流水线的依赖执行方案
以下是几种可行的方案,确保第二个Vertex AI流水线仅在第一个成功完成后启动:
1. 整合为单一流水线的组件依赖
如果两个流水线的逻辑可以整合到同一个工作流中,推荐将它们封装为Kubeflow Pipelines组件,通过显式依赖关系实现顺序执行:
- 将第一个流水线的核心逻辑封装为一个组件,第二个流水线的逻辑封装为另一个组件
- 在定义流水线时,指定第二个组件依赖第一个组件的输出,并通过
after()方法强制执行顺序 - 只有第一个组件成功执行完成,第二个组件才会启动
示例代码:
from kfp.v2 import dsl from kfp.v2.dsl import component # 第一个流水线组件 @component(base_image="python:3.9") def first_pipeline() -> str: # 替换为第一个流水线的实际逻辑 return "first_pipeline_result" # 第二个流水线组件,依赖第一个的输出 @component(base_image="python:3.9") def second_pipeline(input_data: str): # 替换为第二个流水线的实际逻辑,使用第一个的结果 print(f"Received data from first pipeline: {input_data}") # 定义流水线 @dsl.pipeline(name="dependent-pipelines-workflow") def pipeline(): first_task = first_pipeline() # 绑定依赖关系,确保第一个任务成功后执行第二个 second_task = second_pipeline(input_data=first_task.output) second_task.after(first_task)
2. 基于Pub/Sub + Cloud Functions的事件触发
如果两个流水线必须独立部署,可以利用Vertex AI的事件通知机制,结合Cloud Functions实现触发:
- 为第一个流水线配置完成事件的Pub/Sub订阅,当流水线状态变更时,事件会发送到指定主题
- 创建Cloud Functions,订阅该Pub/Sub主题,在函数中解析事件中的流水线状态
- 若第一个流水线状态为
SUCCEEDED,则调用Vertex AI API启动第二个流水线,并传入第一个流水线的输出结果
示例函数逻辑:
import google.cloud.aiplatform as aiplatform def trigger_second_pipeline(event, context): # 从Pub/Sub事件中提取流水线任务名称 pipeline_job_name = event["attributes"]["pipeline_job_name"] # 获取流水线任务详情 pipeline_run = aiplatform.PipelineJob.get(pipeline_job_name) # 检查第一个流水线是否成功 if pipeline_run.state == aiplatform.PipelineJobState.SUCCEEDED: # 提取第一个流水线的输出结果 first_output = pipeline_run.outputs["output_key"].value # 启动第二个流水线 second_pipeline = aiplatform.PipelineJob( display_name="second-pipeline", template_path="gs://your-bucket/pipelines/second_pipeline.json", parameter_values={"input_data": first_output} ) second_pipeline.submit()
3. 使用Cloud Workflows编排流水线
借助Cloud Workflows的流程编排能力,实现两个流水线的顺序执行与状态检查:
- 在Workflows中定义两个步骤:第一步启动第一个流水线,第二步等待其完成并检查状态
- 若第一个流水线成功,则触发第二个流水线;若失败则终止流程
- 可以直接通过Workflows调用Vertex AI的API,无需额外中间组件
示例Workflows配置片段:
main: steps: - run_first_pipeline: call: googleapis.aiplatform.v1.projects.locations.pipelineJobs.run args: parent: projects/your-project/locations/us-central1 pipelineJob: displayName: first-pipeline templateUri: gs://your-bucket/pipelines/first_pipeline.json result: first_job - wait_for_first_completion: call: googleapis.aiplatform.v1.projects.locations.pipelineJobs.get args: name: ${first_job.name} retry: predicate: ${wait_for_first_completion.response.state != "SUCCEEDED" && wait_for_first_completion.response.state != "FAILED"} interval: 60 max_retries: 60 - check_success: switch: - condition: ${wait_for_first_completion.response.state == "SUCCEEDED"} next: run_second_pipeline - condition: ${wait_for_first_completion.response.state == "FAILED"} next: exit_with_failure - run_second_pipeline: call: googleapis.aiplatform.v1.projects.locations.pipelineJobs.run args: parent: projects/your-project/locations/us-central1 pipelineJob: displayName: second-pipeline templateUri: gs://your-bucket/pipelines/second_pipeline.json parameterValues: input_data: ${wait_for_first_completion.response.outputs.output_key.value} - exit_success: return: "Second pipeline started successfully after first pipeline succeeded" - exit_with_failure: raise: "First pipeline failed, second pipeline not initiated"
内容的提问来源于stack exchange,提问作者Shadab Hussain
相关产品推荐
相关产品推荐

