如何配置Dagster每日分区资产依赖前一分区的其他资产?
如何设置Dagster分区资产依赖前一日分区?
要让my_asset_a的T日分区依赖my_asset_b的T-1日分区,适配两者的运行时间差,你可以通过以下两种方式实现:
方法一:使用AssetIn声明式指定分区偏移
这是最简洁的实现方式,直接在资产装饰器中通过ins参数定义依赖资产的分区偏移规则:
from datetime import timedelta from dagster import asset, AssetIn, DailyPartitionsDefinition daily_partitions_def = DailyPartitionsDefinition(start_date="2024-01-01") @asset( partitions_def=daily_partitions_def, ins={"my_asset_b": AssetIn(partition_key_fn=lambda dt: dt - timedelta(days=1))} ) def my_asset_a(context, my_asset_b): dt = context.partition_key return run_a(dt, my_asset_b) @asset(partitions_def=daily_partitions_def) def my_asset_b(context): dt = context.partition_key return run_b(dt)
说明
partition_key_fn会接收当前资产的分区日期(dt),返回依赖资产的目标分区日期,这里通过timedelta(days=1)实现前一日偏移。- 当
my_asset_a运行2024-07-08分区时,会自动加载my_asset_b的2024-07-07分区数据,完全适配两者的运行时间差。
方法二:手动加载指定分区资产
如果需要动态调整依赖逻辑(比如根据条件切换T/T-1依赖),可以在函数内部手动加载目标分区的资产:
from datetime import timedelta from dagster import asset, DailyPartitionsDefinition, AssetKey daily_partitions_def = DailyPartitionsDefinition(start_date="2024-01-01") @asset(partitions_def=daily_partitions_def) def my_asset_a(context): dt = context.partition_key # 计算前一日分区键 prev_dt = (dt - timedelta(days=1)).strftime("%Y-%m-%d") # 手动加载my_asset_b的前一日分区 my_asset_b_prev = context.resources.asset_loader.load_asset( AssetKey("my_asset_b"), partition_key=prev_dt ) return run_a(dt, my_asset_b_prev) @asset(partitions_def=daily_partitions_def) def my_asset_b(context): dt = context.partition_key return run_b(dt)
说明
- 这种方式适合需要灵活调整依赖的场景,比如某些特殊日期需要依赖当日数据时,只需修改
prev_dt的计算逻辑即可。
内容的提问来源于stack exchange,提问作者pyCthon
相关产品推荐
相关产品推荐

