在Dagster资产中获取当前执行日期,是否有更简便的方法?
在Dagster中简化获取当前日期的实现
你当前的写法没问题,但Dagster有更贴合自身生态的简化方案,不用单独定义current_dt依赖函数,以下是几种更高效的方式:
1. 通过资产执行上下文直接拿执行日期
在资产函数中注入AssetExecutionContext,从上下文里直接获取执行时间并格式化:
from dagster import asset, AssetExecutionContext from datetime import datetime @asset def my_task(context: AssetExecutionContext): # 默认是UTC时间,按需转时区 return context.execution_time.strftime('%Y-%m-%d')
如果需要本地时间,改成:
return context.execution_time.astimezone().strftime('%Y-%m-%d')
2. 用分区资产自动传入日期(对标Airflow的ds参数)
如果你的资产是按日调度的,用DailyPartitionsDefinition定义分区后,Dagster会自动把当前分区的日期字符串(格式%Y-%m-%d)通过上下文传给你,和Airflow的ds逻辑完全一致:
from dagster import asset, DailyPartitionsDefinition daily_partitions = DailyPartitionsDefinition(start_date="2024-01-01") @asset(partitions_def=daily_partitions) def my_task(context: AssetExecutionContext): return context.partition_key
这种方式最适合按日运行的定时任务,无需手动处理日期逻辑。
3. 内置配置参数(可选)
也可以通过配置类默认值传递,但灵活性不如前两种:
from dagster import asset, Config from datetime import datetime class TaskConfig(Config): current_dt: str = datetime.today().strftime('%Y-%m-%d') @asset def my_task(config: TaskConfig): return config.current_dt
以上方案里,前两种是Dagster推荐的做法,省去了额外依赖函数的定义,逻辑更简洁,也更符合Dagster的设计思路。
内容的提问来源于stack exchange,提问作者pyCthon
相关产品推荐
相关产品推荐

