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

如何在Dagster中实现仅在所有动态并行Op执行完成后运行指定Op

问题解决:Dagster中等待所有动态并行Op完成后执行后续Op

问题分析

你当前的问题在于,clients.map(...)返回的是DynamicMap对象,直接将其作为Nothing类型的输入传给step4时,Dagster无法识别需要等待所有并行的step3实例执行完成。DynamicMap仅表示一组动态生成的Op,必须通过collect()方法显式告知Dagster等待所有并行任务结束。

解决方案

修改你的Graph定义,将动态映射后的结果通过collect()方法转换为一个等待所有并行Op完成的依赖,再传递给step4:

from dagster import op, graph, DynamicOut, DynamicOutput, In, Nothing

@op()
def step_1_get_response():
    return {'exemple': 'data'}

@op()
def step_2_get_client_list():
    return ['client_1', 'client_2', 'client_3'] #客户数量是动态的

@op(out=DynamicOut())
def parallelize_clients(context, client_list):
    for client in client_list:
        yield DynamicOutput(client, mapping_key=str(client))

@op()
def step_3_update_database_cliente(response, client):
    # 这里填写你的数据库更新逻辑
    print(f"Updated {client} database")

@op(ins={"start": In(Nothing)})
def step_4():
    print("All client updates completed, executing step4")

@graph()
def job_exemple_graph():
    response = step_1_get_response()
    clients_list = step_2_get_client_list()
    clients = parallelize_clients(clients_list)
    
    # 保存动态映射的结果,调用collect()等待所有并行任务完成
    client_updates = clients.map(lambda client: step_3_update_database_cliente(response, client))
    # 将collect()的结果作为step4的触发条件
    step_4(start=client_updates.collect())

关键说明

  • collect()方法会创建一个隐式的汇总Op,它会等待所有由DynamicMap生成的step3实例全部执行完成后,才会向step4发送Nothing信号。
  • 必须将map()的结果赋值给变量,再调用collect(),否则无法建立正确的依赖关系。

内容的提问来源于stack exchange,提问作者Cristiano Cardoso dos Santos

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 21:15:37