曾用Airflow depends_on_past,如何编写依赖自身早期分区的Dagster资产?
在Dagster中实现类似Airflow
depends_on_past的功能 要让每日分区资产依赖自身的早期分区(比如前一天的分区),可以通过以下步骤实现:
1. 定义带自依赖的每日分区资产
首先创建每日分区规则,然后在资产的输入配置里指定依赖自身的前一个分区。核心是用DailyPartitionsDefinition定义分区规则,再通过AssetIn的partition_key_fn动态获取前一天的分区键。
from dagster import ( AssetIn, DailyPartitionsDefinition, asset, build_asset_reconciliation_sensor, Definitions, ) # 定义每日分区规则,起始日期按需调整 daily_partitions = DailyPartitionsDefinition(start_date="2024-01-01") @asset( partitions_def=daily_partitions, # 声明依赖自身的前一个分区 ins={ "previous_self": AssetIn( key="daily_asset", partition_key_fn=lambda partition_key: daily_partitions.get_prev_partition_key(partition_key) ) } ) def daily_asset(context, previous_self): current_partition = context.asset_partition_key_for_output() context.log.info(f"Processing partition {current_partition}") # 处理首次运行的边界情况:第一个分区没有前置依赖 if previous_self is not None: context.log.info(f"Using data from previous partition: {previous_self}") # 替换为你的实际业务逻辑 return {"current_partition": current_partition, "prev_data": previous_self}
代码说明:
DailyPartitionsDefinition:指定资产按日划分分区,起始日期可根据业务需求修改。AssetIn(key="daily_asset"):明确依赖当前资产的另一个分区。partition_key_fn:调用get_prev_partition_key自动计算当前分区的前一日分区键,实现自依赖逻辑。- 资产函数中,
previous_self会传入前一日分区的输出结果,首次运行时该值为None,需自行处理边界场景。
2. 配置资产调和传感器
使用build_asset_reconciliation_sensor监控资产分区状态,当前置分区运行完成后,自动触发当前分区的执行。
# 创建资产调和传感器 daily_asset_sensor = build_asset_reconciliation_sensor( name="daily_asset_reconciliation_sensor", asset_selection=[daily_asset], # 可选:设置传感器检查间隔,比如每小时一次 minimum_interval_seconds=3600, ) # 组装Dagster定义 defs = Definitions( assets=[daily_asset], sensors=[daily_asset_sensor], )
关键提示:
- 确保资产存储(如本地文件、数据库)能区分不同分区的数据,Dagster会自动维护分区的运行元数据。
- 若需要依赖更早的分区(比如前7天),可修改
partition_key_fn逻辑,通过daily_partitions.get_partition_key_offset获取指定偏移量的分区键。
内容的提问来源于stack exchange,提问作者Sandy Ryza
相关产品推荐
相关产品推荐

