如何在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
相关产品推荐
相关产品推荐

