Dagster中为同一资产依赖配置多个TimeWindowPartitionMapping的方法
问题解决:同一资产配置多分区映射及KeyError修复
可以为同一资产配置不同的TimeWindowPartitionMapping,你的KeyError是代码中的几个错误导致的,以下是问题分析和修复方案:
核心错误点
- 资产引用错误:你的
SourceAssetkey是merged_data,但AssetIn里写的是original_dfs,Dagster无法找到对应资产。 - 参数名不匹配:
ins中previous_day_original_dfs末尾多了空格,且函数参数名是previous_original_dfs,两者名称不一致导致映射失败。 - 时间偏移计算错误: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
相关产品推荐
相关产品推荐

