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

在Dagster中使用mem_io_manager进程内执行作业后如何获取资产值?

解决方案:从进程内执行结果中提取资产值

当使用mem_io_manager时,由于它仅将数据存储在当前进程的内存中,作业执行完成后无法通过load_asset_value直接读取(这也是你遇到报错的核心原因)。但可以直接从execute_in_process返回的结果对象中提取资产输出值,无需依赖IO管理器的后续读取操作。

方法1:简洁获取单个资产值

利用execute_in_process返回的ExecuteInProcessResult对象,通过资产对应的节点名称直接获取输出(资产的节点名默认与资产名一致):

from dagster import asset, repository

@asset
def my_int():
    return 1

@repository
def my_repo():
    return [my_int]

# 执行资产作业
result = my_repo.get_job('__ASSET_JOB').execute_in_process()

# 直接获取指定资产的输出值
my_int_value = result.output_for_node("my_int")
print(my_int_value)  # 输出: 1

方法2:实现批量资产值获取(类似你设想的return_assets)

封装工具函数,支持批量传入资产对象或资产名,一次性返回多个资产的值:

from dagster import AssetKey

def get_assets_from_execution(result, assets):
    asset_values = {}
    for item in assets:
        # 兼容传入资产对象或资产名字符串
        if isinstance(item, str):
            asset_key = AssetKey(item)
            node_name = item
        else:
            asset_key = item.key
            node_name = asset_key.path[-1]
        
        asset_values[asset_key] = result.output_for_node(node_name)
    return asset_values

# 使用示例
result = my_repo.get_job('__ASSET_JOB').execute_in_process()
my_assets = get_assets_from_execution(result, [my_int])
print(my_assets[my_int.key])  # 输出: 1

补充说明:为什么load_asset_value对mem_io_manager无效?

load_asset_value的设计是通过仓库配置的IO管理器读取持久化的资产数据,但mem_io_manager仅在作业执行的进程内临时存储数据,执行结束后内存中的数据会被释放。此外它的实现依赖执行时的上下文参数(如step_key),离线调用时无法提供这些参数,因此会抛出你看到的错误。而直接从执行结果中获取值,是因为结果对象保留了执行过程中的输出快照,无需再通过IO管理器读取。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 05:46:09