如何按group_name筛选Dagster资产并创建对应分组作业?
按指定group_name创建独立Dagster作业
当然可以仅加载指定分组的资产,不用修改资产的加载逻辑,直接在作业的selection参数里用Dagster的资产选择语法即可,给你两种常用实现方式:
1. 使用资产选择字符串
直接通过"group:分组名"的语法筛选对应分组的资产,代码示例:
from dagster import define_asset_job from ..assets import my_assets # 创建group1专属作业 my_group1_job = define_asset_job( name="group1_job", selection="group:group1", # 核心:用group筛选语法指定分组 description="仅加载group1的资产数据" ) # 创建group2专属作业 my_group2_job = define_asset_job( name="group2_job", selection="group:group2", description="仅加载group2的资产数据" )
2. 使用AssetSelection类(更推荐,逻辑更清晰)
导入AssetSelection类,通过其groups方法指定目标分组,适合后续扩展复杂的资产选择逻辑:
from dagster import define_asset_job, AssetSelection from ..assets import my_assets my_group1_job = define_asset_job( name="group1_job", selection=AssetSelection.groups("group1"), description="仅加载group1的资产数据" ) my_group2_job = define_asset_job( name="group2_job", selection=AssetSelection.groups("group2"), description="仅加载group2的资产数据" )
备选方案:先加载所有资产再筛选
如果需要先加载全部资产再做额外处理,可使用filter_assets函数过滤指定分组:
from dagster import define_asset_job, filter_assets, load_assets_from_modules from ..assets import my_assets # 加载所有资产 all_assets = load_assets_from_modules([my_assets]) # 筛选出group1的资产 group1_assets = filter_assets(all_assets, group_name="group1") my_group1_job = define_asset_job( name="group1_job", selection=group1_assets, description="仅加载group1的资产数据" )
注意:前两种方法更高效,Dagster会在作业定义阶段直接筛选资产,无需加载全部资产后再过滤。
内容的提问来源于stack exchange,提问作者x89
相关产品推荐
相关产品推荐

