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

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)

修改说明

  1. 明确依赖关系:每个clean_results_dir直接依赖对应模型Op的输出,确保模型运行完成后立即触发清理,解决顺序混乱问题。
  2. 聚合清理结果:使用collect函数汇总所有清理Op的输出,让generate_reports等待所有清理操作完成后再执行,满足“所有模型运行完毕(无论成败)再生成报表”的要求。
  3. 独立执行单元:各模型-清理单元之间无依赖,配合max_concurrent:1的设置,会串行执行每个单元,同时某个模型失败时,后续单元仍能正常运行。

内容的提问来源于stack exchange,提问作者Rafa Calvo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 03:15:55