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

Dagster中为同一资产依赖配置多个TimeWindowPartitionMapping的方法

问题解决:同一资产配置多分区映射及KeyError修复

可以为同一资产配置不同的TimeWindowPartitionMapping,你的KeyError是代码中的几个错误导致的,以下是问题分析和修复方案:

核心错误点

  1. 资产引用错误:你的SourceAsset key是merged_data,但AssetIn里写的是original_dfs,Dagster无法找到对应资产。
  2. 参数名不匹配:ins中previous_day_original_dfs末尾多了空格,且函数参数名是previous_original_dfs,两者名称不一致导致映射失败。
  3. 时间偏移计算错误:15分钟分区一天对应96个间隔(24*4),你写的-4*24*7是一周的偏移,不符合前一天的需求。

修复后的完整代码

original_dfs = SourceAsset(
    key="merged_data",
    io_manager_key="partitioned_parquet_io_manager",
    description="Descriptions",
    partitions_def=fifteen_minutes_partitions,
    metadata={"io_manager": TimePartitionConfig(time_partition_column="timestamp")},
)


@asset(
    io_manager_key="partitioned_parquet_io_manager",
    partitions_def=fifteen_minutes_partitions,
    description="description",
    metadata={"io_manager": TimePartitionConfig(time_partition_column="timestamp")},
    ins={
        "previous_day_original_dfs": AssetIn(
            "merged_data",  # 修复:引用正确的资产key
            partition_mapping=TimeWindowPartitionMapping(start_offset=-96, end_offset=-96),  # 修复:15分钟分区一天对应96个偏移
        ),
        "current_original_dfs": AssetIn(
            "merged_data",  # 修复:引用正确的资产key
            partition_mapping=TimeWindowPartitionMapping(start_offset=0, end_offset=0),
        ),
    },
)
def lagged_df(
    context: OpExecutionContext,
    previous_day_original_dfs: pd.DataFrame,  # 修复:参数名与ins中的key一致
    current_original_dfs: pd.DataFrame,
) -> pd.DataFrame:
    # 此处添加你的业务逻辑
    return pd.concat([previous_day_original_dfs, current_original_dfs], axis=1)

替代方案:手动加载指定分区

如果因特殊场景无法通过多输入实现,可在资产内部通过context手动加载目标分区数据:

@asset(
    io_manager_key="partitioned_parquet_io_manager",
    partitions_def=fifteen_minutes_partitions,
    description="description",
    metadata={"io_manager": TimePartitionConfig(time_partition_column="timestamp")},
    ins={"original_dfs": AssetIn("merged_data")},
)
def lagged_df(
    context: OpExecutionContext,
    original_dfs: pd.DataFrame,
    io_manager: IOManager,
) -> pd.DataFrame:
    # 获取当前分区并计算前一天分区key(根据你的分区格式调整)
    from datetime import datetime, timedelta
    current_dt = datetime.fromisoformat(context.partition_key)
    previous_day_dt = current_dt - timedelta(days=1)
    previous_partition_key = previous_day_dt.isoformat()
    
    # 手动加载前一天分区数据
    previous_day_data = io_manager.load_input(
        context=context,
        input_def=InputDefinition(name="previous_day_data"),
        asset_key=AssetKey("merged_data"),
        partition_key=previous_partition_key,
    )
    
    # 业务逻辑处理
    return pd.concat([previous_day_data, original_dfs], axis=1)

内容的提问来源于stack exchange,提问作者Loïc Quivron

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:50:00