Dagster使用I/O Manager加载源数据报错,求排查方案
Dagster输入资产未找到错误的解决方案
问题场景
刚接触Dagster,尝试通过自定义I/O Manager读取CSV文件转为pandas DataFrame,再写入Parquet文件。编写的资产代码运行时触发以下错误:
DagsterInvalidDefinitionError: Input asset '["csv_asset"]' for asset '["out", "parquet_asset"]' is not produced by any of the provided asset ops and is not one of the provided sources.
用户原始代码:
import pandas as pd from dagster import asset, AssetKey, AssetSpec csv_asset = AssetSpec( key = AssetKey(["in", "dropout"]), metadata = {"dagster/io_manager_key": "local_csv_io_manager"}, ) @asset(io_manager_key="local_parquet_io_manager", key_prefix=["out"]) def parquet_asset(csv_asset: pd.DataFrame): return csv_asset defs = Definitions( assets=[csv_asset, parquet_asset], resources={ "local_parquet_io_manager": LocalParquetIOManager(), "local_csv_io_manager": LocalCSVIOManager() } )
错误原因
Dagster默认会将资产函数的参数名解析为资产键(AssetKey)。这里parquet_asset的参数名是csv_asset,所以Dagster会尝试寻找键为["csv_asset"]的资产,但实际定义的源资产键是["in", "dropout"],导致匹配失败。
解决方案
1. 明确指定输入资产的键
使用@asset装饰器的ins参数,通过AssetIn关联正确的源资产键:
import pandas as pd import os from dagster import asset, AssetKey, AssetSpec, Definitions, AssetIn, IOManager, InputContext, OutputContext csv_asset = AssetSpec( key = AssetKey(["in", "dropout"]), metadata = {"dagster/io_manager_key": "local_csv_io_manager"}, ) @asset( io_manager_key="local_parquet_io_manager", key_prefix=["out"], ins={"csv_asset": AssetIn(key=AssetKey(["in", "dropout"]))} # 明确关联源资产 ) def parquet_asset(csv_asset: pd.DataFrame): return csv_asset # 实现自定义CSV I/O Manager class LocalCSVIOManager(IOManager): def __init__(self, base_dir: str = "data/csv"): self.base_dir = base_dir os.makedirs(base_dir, exist_ok=True) def load_input(self, context: InputContext) -> pd.DataFrame: # 根据资产键拼接CSV文件路径 file_name = f"{context.asset_key.path[-1]}.csv" file_path = os.path.join(self.base_dir, file_name) return pd.read_csv(file_path) # 实现自定义Parquet I/O Manager class LocalParquetIOManager(IOManager): def __init__(self, base_dir: str = "data/parquet"): self.base_dir = base_dir os.makedirs(base_dir, exist_ok=True) def handle_output(self, context: OutputContext, obj: pd.DataFrame): file_name = f"{context.asset_key.path[-1]}.parquet" file_path = os.path.join(self.base_dir, file_name) obj.to_parquet(file_path) defs = Definitions( assets=[csv_asset, parquet_asset], resources={ "local_parquet_io_manager": LocalParquetIOManager(), "local_csv_io_manager": LocalCSVIOManager() } )
2. 使用命名空间风格的参数名
Dagster支持用双下划线分隔资产键的命名空间层级,直接用in__dropout作为参数名,对应AssetKey(["in", "dropout"]):
@asset(io_manager_key="local_parquet_io_manager", key_prefix=["out"]) def parquet_asset(in__dropout: pd.DataFrame): return in__dropout
3. 确保I/O Manager正确实现
必须完成核心方法:
LocalCSVIOManager:实现load_input方法,读取指定路径的CSV文件并返回DataFrameLocalParquetIOManager:实现handle_output方法,将DataFrame写入Parquet文件
关键注意点
- 导入缺失的模块:确保代码中导入了
Definitions和AssetIn(使用第一种方案时) - 源资产(
AssetSpec定义的资产)不需要实现handle_output,因为它是读取外部文件的源,而非生成数据的资产
内容的提问来源于stack exchange,提问作者jbelis
相关产品推荐
相关产品推荐

