如何让Kubeflow Pipeline最后组件不受前置分支状态影响执行?
Kubeflow Pipeline 条件分支后必执行组件的解决方案
问题描述
我搭建的Kubeflow Pipeline以c0_load_config为起点,分为变量和事件两个分支:
- 变量分支:执行
c1_var_load后,通过kfp.dsl.Condition判断是否执行后续的c2_var_procesado和c3_var_escritura组件; - 事件分支:执行
c1_eve_load_geo、c2_eve_procesado后,通过kfp.dsl.Condition判断是否执行c3_eve_escritura组件。
两个分支处理完成后,需要执行公共组件c4_salud_unificado——它不接收前序分支的参数,但必须等两个分支都处理完成(无论分支内的Condition是否触发)才执行。但当前代码中使用.after(cond_variables).after(cond_eventos)的方式,只要任一Condition分支被跳过,c4_salud_unificado就不会触发。
原代码如下:
@kfp.dsl.pipeline(name='name') def pipeline_hvac(some_parameters:str): # Pipeline Varaibles # c0_load_config_task = c0_load_config(some_parameters=some_parameters) c1_var_load_task = c1_var_load(config_path=c0_load_config_task.outputs['config_path']) with kfp.dsl.Condition(condition=(c1_var_load_task.outputs['output'] == 'true')) as cond_variables: c2_var_process_task = c2_var_procesado( config_path=c0_load_config_task.outputs['config_path'], df_var_diario_path=c1_var_load_task.outputs['df_var_clean_path']) c3_var_escritura( config_path=c0_load_config_task.outputs['config_path'], df_variables_estado_dia_path=c2_var_process_task.outputs['df_variables_estado_dia_path'], df_estado_eaa_path=c2_var_process_task.outputs['df_estado_eaa_path']) c1_eve_load_geo_task = c1_eve_load_geo( config_path=c0_load_config_task.outputs['config_path']) c2_eve_procesado_task = c2_eve_procesado(config_path=c0_load_config_task.outputs['config_path']).after(c1_eve_load_geo_task) with kfp.dsl.Condition(condition=(c2_eve_procesado_task.outputs['output'] == 'true')) as cond_eventos: c3_eve_escritura( config_path=c0_load_config_task.outputs['config_path'], dict_results_path=c2_eve_procesado_task.outputs['dict_results_path']) c4_salud_unificado( config_path = c0_load_config_task.outputs['config_path']).after(cond_variables).after(cond_eventos)
解决方案
问题原因
直接依赖Condition对象时,只有当Condition分支被实际执行(进入分支内部运行组件),这个依赖才会被标记为完成;如果Condition不满足导致分支被跳过,Condition对象对应的节点不会进入完成状态,进而阻塞后续组件的执行。
修改方法
让c4_salud_unificado依赖每个分支一定会执行完成的最后一个任务,而非Condition对象:
- 变量分支中,
c1_var_load_task是一定会执行的(无论后续Condition是否触发) - 事件分支中,
c2_eve_procesado_task是一定会执行的(无论后续Condition是否触发)
修改后的完整代码:
@kfp.dsl.pipeline(name='name') def pipeline_hvac(some_parameters:str): # Pipeline Varaibles # c0_load_config_task = c0_load_config(some_parameters=some_parameters) c1_var_load_task = c1_var_load(config_path=c0_load_config_task.outputs['config_path']) with kfp.dsl.Condition(condition=(c1_var_load_task.outputs['output'] == 'true')) as cond_variables: c2_var_process_task = c2_var_procesado( config_path=c0_load_config_task.outputs['config_path'], df_var_diario_path=c1_var_load_task.outputs['df_var_clean_path']) c3_var_escritura( config_path=c0_load_config_task.outputs['config_path'], df_variables_estado_dia_path=c2_var_process_task.outputs['df_variables_estado_dia_path'], df_estado_eaa_path=c2_var_process_task.outputs['df_estado_eaa_path']) c1_eve_load_geo_task = c1_eve_load_geo( config_path=c0_load_config_task.outputs['config_path']) c2_eve_procesado_task = c2_eve_procesado(config_path=c0_load_config_task.outputs['config_path']).after(c1_eve_load_geo_task) with kfp.dsl.Condition(condition=(c2_eve_procesado_task.outputs['output'] == 'true')) as cond_eventos: c3_eve_escritura( config_path=c0_load_config_task.outputs['config_path'], dict_results_path=c2_eve_procesado_task.outputs['dict_results_path']) # 修改:依赖两个分支中一定会执行的最终任务,而非Condition对象 c4_salud_unificado( config_path = c0_load_config_task.outputs['config_path'] ).after(c1_var_load_task).after(c2_eve_procesado_task)
这样修改后,不管两个分支里的Condition是否触发,只要c1_var_load_task和c2_eve_procesado_task都执行完成,c4_salud_unificado就会自动触发执行。
内容的提问来源于stack exchange,提问作者Alejandro de la Parra
相关产品推荐
相关产品推荐

