如何通过Python API在Dagster中独立物化下游资产?
我使用dagster==1.0.11,希望通过Python API独立物化下游资产downstream_asset,无需依赖上游资产upstream_asset的元数据,只要上游资产的持久化结果存在即可成功执行。
示例代码
# example.py from dagster import asset, materialize, repository @asset def upstream_asset(): return [1, 2, 3] @asset def downstream_asset(upstream_asset): return upstream_asset + [4] @repository def repo(): return [upstream_asset, downstream_asset] if __name__ == "__main__": materialize([downstream_asset])
执行错误
执行python example.py时抛出错误:
dagster._core.errors.DagsterInvalidDefinitionError: Input asset '["upstream_asset"]' for asset '["downstream_asset"]' is not produced by any of the provided asset ops and is not one of the provided sources
完整错误信息:
dagster._core.errors.DagsterExecutionLoadInputError: Error occurred while loading input "upstream_asset" of step "downstream_asset": File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/execution/plan/execute_plan.py", line 224, in dagster_event_sequence_for_step for step_event in check.generator(step_events): File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/execution/plan/execute_step.py", line 320, in core_dagster_event_sequence_for_step step_input.source.load_input_object(step_context, input_def) File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/execution/plan/inputs.py", line 201, in load_input_object yield from _load_input_with_input_manager(loader, load_input_context) File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/execution/plan/inputs.py", line 867, in _load_input_with_input_manager value = input_manager.load_input(context) File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/contextlib.py", line 99, in __exit__ self.gen.throw(type, value, traceback) File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/execution/plan/utils.py", line 82, in solid_execution_error_boundary ) from e The above exception was caused by the following exception: FileNotFoundError: [Errno 2] No such file or directory: '/Users/pedro.viana/dev/nu/mock-model/dagster_home/storage/upstream_asset' File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/execution/plan/utils.py", line 47, in solid_execution_error_boundary yield File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/execution/plan/inputs.py", line 867, in _load_input_with_input_manager value = input_manager.load_input(context) File "/Users/pedro.viana/miniforge3/envs/dagtest/lib/python3.6/site-packages/dagster/_core/storage/fs_io_manager.py", line 181, in load_input with open(filepath, self.read_mode) as read_obj:
期望行为
实现与dagit UI相同的逻辑:先物化upstream_asset成功后,删除DAGSTER_HOME中除storage/upstream_asset外的所有内容(仅保留上游资产的持久化结果),重启dagit后选择downstream_asset点击“Materialize selected”即可成功执行——此时系统仅检查IO Manager指定路径下的文件是否存在,无需上游资产的元数据。
背景
有多套运行相同代码版本(使用s3_io_manager)的Dagster实例,若某实例物化了上游资产,其他实例执行下游资产时应能直接使用该持久化结果。
要实现该需求,需将上游资产标记为外部资产(external asset),告知Dagster该资产不由当前执行流程生成,而是依赖外部已存在的持久化结果。具体修改如下:
修改后的代码
# example.py from dagster import asset, materialize, repository, AssetKey, SourceAsset # 将上游资产定义为外部资产 upstream_asset = SourceAsset(key=AssetKey("upstream_asset")) @asset def downstream_asset(upstream_asset): return upstream_asset + [4] @repository def repo(): return [upstream_asset, downstream_asset] if __name__ == "__main__": materialize([downstream_asset])
关键说明
SourceAsset的作用:通过SourceAsset标记后,Dagster会认为upstream_asset的数据源来自外部存储(本地文件、S3等),不再要求当前流程中包含该资产的生成逻辑。- 执行逻辑变化:调用
materialize([downstream_asset])时,Dagster会直接通过IO Manager尝试加载upstream_asset的持久化结果,只要对应存储路径存在数据,就能成功执行。 - 多实例适配:使用
s3_io_manager时,只要所有实例共享同一S3存储桶,任意实例物化的上游资产结果,其他实例都能直接加载使用,无需同步元数据。
验证步骤
- 先通过任意Dagster实例物化
upstream_asset(可使用原代码执行materialize([upstream_asset]),或在dagit中操作)。 - 替换为上述修改后的代码,执行
python example.py,即可直接加载已存在的上游资产结果,成功物化downstream_asset。
内容的提问来源于stack exchange,提问作者pedrovgp

