如何基于不同分区与调度的相似资产创建多个Dagster Job
为单个资产创建独立Job并配置不同调度的实现方式
无需拆分assets.py文件,你可以直接通过指定资产选择器或直接引用资产对象的方式,为每个资产单独创建Job并配置不同调度,具体实现如下:
1. 主仓库文件(如repo.py)的实现代码
导入依赖
from dagster import ( repository, define_asset_job, ScheduleDefinition, AssetSelection, load_assets_from_modules ) # 导入assets.py中的资产对象 from assets import get_rfr_final_asset, get_rfr_provisional_asset # 或者直接导入整个assets模块 import assets
方式一:直接引用资产对象创建Job
适合资产数量较少的场景,直接指定单个资产:
# 为get_rfr_final_asset创建独立Job final_rfr_job = define_asset_job( name="final_rfr_job", selection=[get_rfr_final_asset], ) # 为get_rfr_provisional_asset创建独立Job provisional_rfr_job = define_asset_job( name="provisional_rfr_job", selection=[get_rfr_provisional_asset], )
方式二:通过AssetSelection按资产名称筛选
适合资产数量较多的场景,先加载所有资产再精准筛选:
# 加载assets模块中的所有资产 all_assets = load_assets_from_modules([assets]) # 按资产key筛选创建final资产的Job final_rfr_job = define_asset_job( name="final_rfr_job", selection=AssetSelection.keys("get_rfr_final_asset"), ) # 按资产key筛选创建provisional资产的Job provisional_rfr_job = define_asset_job( name="provisional_rfr_job", selection=AssetSelection.keys("get_rfr_provisional_asset"), )
为两个Job配置不同调度
根据业务需求设置不同的cron表达式和时区:
# final资产调度:每天凌晨1点执行 final_rfr_schedule = ScheduleDefinition( job=final_rfr_job, cron_schedule="0 1 * * *", execution_timezone="Asia/Shanghai", # 按需指定时区 ) # provisional资产调度:每天凌晨3点执行 provisional_rfr_schedule = ScheduleDefinition( job=provisional_rfr_job, cron_schedule="0 3 * * *", execution_timezone="Asia/Shanghai", )
注册仓库
将Job和调度注册到Dagster仓库:
@repository def my_repository(): return [ final_rfr_job, provisional_rfr_job, final_rfr_schedule, provisional_rfr_schedule, ]
2. 核心逻辑说明
- 通过
define_asset_job的selection参数精准指定单个资产,避免所有资产归为一组 - 每个Job可独立配置调度规则,完全隔离两个资产的运行周期
- 无需修改原
assets.py文件,保持代码结构简洁
内容的提问来源于stack exchange,提问作者pyCthon
相关产品推荐
相关产品推荐

