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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 06:05:14