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

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文件并返回DataFrame
  • LocalParquetIOManager:实现handle_output方法,将DataFrame写入Parquet文件

关键注意点

  • 导入缺失的模块:确保代码中导入了Definitions和AssetIn(使用第一种方案时)
  • 源资产(AssetSpec定义的资产)不需要实现handle_output,因为它是读取外部文件的源,而非生成数据的资产

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 01:04:54