如何在Dagster中按调度运行指定分区的多分区资产任务
问题描述
- 多分区配置定义:
multiple_partition = MultiPartitionsDefinition({ "date": DailyPartitionsDefinition(start_date="2024-06-06", end_offset=1), "cycle": StaticPartitionsDefinition(['u0', 'u1', 'u2', 'u3', 'u4']), })
- 资产依赖链:
extract_AUL -> transform_AUL、extract_FV -> transform_FV、extract_TN -> transform_TN - 需求:创建5组独立的调度+任务,每组仅运行上述6个资产的指定cycle分区(如
schedule_u0+job_u0对应cycle='u0'),但当前配置会触发所有资产的所有分区,尝试通过ops_config指定分区无效。
可行实现方案
核心思路:通过Partition Selection为每个任务固定指定的cycle分区,同时保留date分区的动态调度逻辑。
1. 修改任务创建函数
使用PartitionSelection过滤仅包含指定cycle的分区,date分区自动取最新值:
from dagster import define_asset_job, PartitionSelection, MultiPartitionsDefinition def create_asset_job(cycle: str, partition: MultiPartitionsDefinition, assets: list): # 固定cycle分区,date分区取最新可用值 partition_selection = PartitionSelection.fixed( {"cycle": cycle, "date": PartitionSelection.latest()} ) return define_asset_job( name=f"job_{cycle}", selection=assets, partitions_def=partition, partition_selection=partition_selection )
2. 简化调度创建函数
直接传入对应cycle的任务即可,无需额外配置:
from dagster import build_schedule_from_partitioned_job def create_schedule(cycle: str, hour: int, minute: int, partition: MultiPartitionsDefinition, assets: list): target_job = create_asset_job(cycle, partition, assets) return build_schedule_from_partitioned_job( job=target_job, description=f"调度周期: {cycle}", name=f"schedule_{cycle}", minute_of_hour=minute, hour_of_day=hour )
3. 批量生成所有调度与任务
遍历所有cycle值,快速创建对应任务和调度:
# 假设assets为包含6个目标资产的列表 # assets = [extract_AUL, transform_AUL, extract_FV, transform_FV, extract_TN, transform_TN] cycles = ['u0', 'u1', 'u2', 'u3', 'u4'] jobs = [] schedules = [] for idx, cycle in enumerate(cycles): # 可根据需求为每个cycle设置不同调度时间 job = create_asset_job(cycle, multiple_partition, assets) schedule = create_schedule(cycle, hour=8+idx, minute=0, partition=multiple_partition, assets=assets) jobs.append(job) schedules.append(schedule)
关键说明
- 之前的
ops_config方式无效,因为多分区规则绑定在资产/任务的partitions_def上,必须通过partition_selection过滤分区,而非通过op config传递参数。 PartitionSelection.fixed确保每个任务仅处理指定cycle的最新date分区,符合每日调度的业务逻辑。- 每组任务和调度独立对应一个cycle,运行时只会触发该cycle下的资产分区,不会执行其他cycle的内容。
内容的提问来源于stack exchange,提问作者Hing Liu
相关产品推荐
相关产品推荐

