You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 11:36:13