Dagster作业问题:模型运行后即时清理、最终生成报告实现困惑
解决Dagster中模型运行与清理步骤顺序混乱的问题
需求目标
- 通过Op工厂创建带指定参数的模型运行Op
- 模型运行失败时后续模型仍需运行(通过设置最大并发数为1避免上下游依赖)
- 每个模型运行完成后立即调用
clean_results_dir清理输出目录 - 所有模型运行完毕后(无论是否失败)才执行
generate_reports方法
遇到的问题
实现第3、4点时出现异常:作业运行时,model1、model2会先全部执行完毕,之后才依次执行cleandir1、cleandir2,顺序不符合“每个模型完成后立即清理”的预期。
原代码
@graph def run_all_models_graph(): model_configs = [{"model_name":"my/model/url/","numPaths":10},{"model_name_two":"my/model/url_two/","numPaths":33}] last_op = None for config in model_configs: model_name = config["model_name"] model_name_formatted = model_name.split("/")[-1] model_instance = get_from_op_factory(model_name_formatted, model_name, config["numPaths"]) last_op = clean_results_dir(model_instance()) return last_op @job( config={ "execution":{ "config": { "multiprocess": { "max_concurrent":1 } } } } ) def run_all_models_job(): result = run_all_models_graph() generate_reports(start_after=result)
问题原因
原循环中,每次仅将最后一个清理Op赋值给last_op,导致Dagster无法识别每个清理Op与对应模型Op的一对一依赖关系,会先调度所有模型Op执行,待全部模型跑完后再执行清理Op。同时,generate_reports仅依赖最后一个清理Op的结果,无法保证等待所有清理完成。
解决方案
修改graph逻辑,明确每个模型与清理的串行依赖,并收集所有清理Op的结果,确保generate_reports等待所有清理完成:
from dagster import graph, job, collect @graph def run_all_models_graph(): model_configs = [{"model_name":"my/model/url/","numPaths":10},{"model_name":"my/model/url_two/","numPaths":33}] clean_ops = [] for config in model_configs: model_name = config["model_name"] model_name_formatted = model_name.split("/")[-1] # 创建模型运行Op model_run_op = get_from_op_factory(model_name_formatted, model_name, config["numPaths"]) # 绑定模型Op与对应清理Op的依赖:模型完成后立即执行清理 clean_op = clean_results_dir(model_run_op()) clean_ops.append(clean_op) # 收集所有清理Op的结果,确保后续报表生成等待所有清理完成 return collect(clean_ops) @job( config={ "execution":{ "config": { "multiprocess": { "max_concurrent":1 } } } } ) def run_all_models_job(): all_clean_results = run_all_models_graph() generate_reports(start_after=all_clean_results)
修改说明
- 明确依赖关系:每个
clean_results_dir直接依赖对应模型Op的输出,确保模型运行完成后立即触发清理,解决顺序混乱问题。 - 聚合清理结果:使用
collect函数汇总所有清理Op的输出,让generate_reports等待所有清理操作完成后再执行,满足“所有模型运行完毕(无论成败)再生成报表”的要求。 - 独立执行单元:各模型-清理单元之间无依赖,配合
max_concurrent:1的设置,会串行执行每个单元,同时某个模型失败时,后续单元仍能正常运行。
内容的提问来源于stack exchange,提问作者Rafa Calvo
相关产品推荐
相关产品推荐

