如何指定Dagster资产的物化存储路径?
如何指定Dagster资产的存储位置
Dagster中资产的存储位置由IOManager控制,以下是几种实用的配置方式:
1. 全局配置默认存储路径
使用内置的FileSystemIOManager,通过base_dir参数指定统一的根目录,所有资产会按{base_dir}/{asset_key}的结构存储:
from dagster import asset, Definitions, FilesystemIOManager @asset def sample_asset(): return {"key": "value"} defs = Definitions( assets=[sample_asset], resources={ "io_manager": FilesystemIOManager(base_dir="/your/target/storage/path") } )
2. 为单个资产指定独立存储路径
通过定义多个IOManager资源,用io_manager_key为不同资产绑定不同的存储配置:
from dagster import asset, Definitions, FilesystemIOManager # 定义两个不同路径的IOManager user_data_io = FilesystemIOManager(base_dir="/storage/user_data") log_data_io = FilesystemIOManager(base_dir="/storage/log_archives") @asset(io_manager_key="user_data_io") def user_profile_asset(): return {"user_id": 1, "name": "Alice"} @asset(io_manager_key="log_data_io") def app_log_asset(): return ["2024-05-20: user login", "2024-05-20: data updated"] defs = Definitions( assets=[user_profile_asset, app_log_asset], resources={ "user_data_io": user_data_io, "log_data_io": log_data_io } )
3. 自定义IOManager实现灵活存储逻辑
如果内置IOManager无法满足需求(比如特殊序列化格式、云存储适配),可以自定义IOManager类来完全控制存储路径和读写逻辑:
from dagster import IOManager, asset, Definitions import os import json class CustomStorageIOManager(IOManager): def __init__(self, root_path: str): self.root_path = root_path def handle_output(self, context, obj): # 根据资产键构造唯一存储路径 asset_path = os.path.join(self.root_path, *context.asset_key.path) os.makedirs(os.path.dirname(asset_path), exist_ok=True) # 用JSON格式序列化数据 with open(f"{asset_path}.json", "w") as f: json.dump(obj, f) def load_input(self, context): asset_path = os.path.join(self.root_path, *context.asset_key.path) with open(f"{asset_path}.json", "r") as f: return json.load(f) # 使用自定义IOManager defs = Definitions( assets=[sample_asset], resources={ "io_manager": CustomStorageIOManager(root_path="/custom/storage/root") } )
内容的提问来源于stack exchange,提问作者MYK
相关产品推荐
相关产品推荐

