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

