如何在Dagster中将数据传递至不同模块的Op?
跨Dagster Job传递数据的可行方案
你当前的写法不可行,因为Dagster的Job是独立的执行单元,默认无法直接共享运行时的Op输出,而且job2里的decrease_num调用时也没有传入所需的参数。下面是几种实现数据跨Job传递的常用方式:
1. 用资产(Assets)替代独立Job
资产是Dagster管理数据依赖的核心概念,适合在不同Job间共享数据。可以把generate_num定义为资产,让两个Job都引用它:
# assets.py from dagster import asset, job, op @asset def generate_num_asset(): return 3 # job1.py from assets import generate_num_asset @op() def increase_num(num): return num + 1 @job() def increment_up(): increase_num(generate_num_asset()) # job2.py from assets import generate_num_asset @op() def decrease_num(num): return num - 1 @op() def multiple_num(num): return num * 2 @job() def get_multiple(): multiple_num(decrease_num(generate_num_asset()))
两个Job会自动依赖generate_num_asset,Dagster会处理数据的生成与复用,你可以通过资产目录查看完整的数据流转链路。
2. 用IO管理器持久化数据
如果必须保留两个独立Job,可以通过IO管理器把generate_num的输出持久化到本地文件、数据库等存储,再在job2中读取:
步骤1:配置本地文件IO管理器
# io_config.py from dagster import io_manager, fs_io_manager @io_manager def local_fs_io_manager(): return fs_io_manager(base_dir="./dagster_data")
步骤2:修改job1,持久化generate_num的输出
# job1.py from dagster import op, job, Out from io_config import local_fs_io_manager @op(out=Out(io_manager_key="local_fs")) def generate_num(): return 3 @op() def increase_num(generate_num): return generate_num + 1 @job(resource_defs={"local_fs": local_fs_io_manager}) def increment_up(): increase_num(generate_num())
步骤3:修改job2,读取持久化的数据
# job2.py from dagster import op, job, In from io_config import local_fs_io_manager @op(in=In(io_manager_key="local_fs")) def decrease_num(generate_num): return generate_num - 1 @op() def multiple_num(decrease_num): return decrease_num * 2 @job(resource_defs={"local_fs": local_fs_io_manager}) def get_multiple(): multiple_num(decrease_num())
注意:这种方式需要先运行increment_up生成并持久化数据,再运行get_multiple读取,要保证执行顺序。
3. 合并为单个Job(业务逻辑允许时)
如果不需要严格拆分Job,可以直接把所有Op整合到一个Job里,直接传递数据:
from dagster import op, job @op() def generate_num(): return 3 @op() def increase_num(generate_num): return generate_num + 1 @op() def decrease_num(generate_num): return generate_num - 1 @op() def multiple_num(decrease_num): return decrease_num * 2 @job() def combined_job(): num = generate_num() increase_num(num) multiple_num(decrease_num(num))
内容的提问来源于stack exchange,提问作者walkingwithataco
相关产品推荐
相关产品推荐

