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

如何用Dagster创建带国家参数的可复用资产并分国调度

在Dagster中实现单文件可复用的多国家资产与调度

核心思路

通过工厂函数批量生成不同国家的资产,结合配置类传递country参数,同时为每个国家的资产组单独配置调度,彻底避免重复编写代码。

具体实现步骤

1. 定义基础配置与业务类

先定义资产所需的配置类,以及你的业务逻辑类(假设DataLoader、RunMyAsset等业务类已存在,这里补充示例框架):

from dagster import asset, DailyPartitionsDefinition, AssetExecutionContext, Config

# 通用日度分区定义
daily_partitions_def = DailyPartitionsDefinition(start_date="2024-01-01")

# 资产配置类:用于传递country参数
class CountryAssetConfig(Config):
    country: str

# 模拟你的业务实现类
class DataLoader:
    def __init__(self, country: str):
        self.country = country

class RunMyAsset:
    def __init__(self, data_loader: DataLoader, country: str):
        self.data_loader = data_loader
        self.country = country
    
    def calc(self, dt: str):
        return f"[{self.country}] 核心资产数据: {dt}"

class RunMyAsset2:
    def __init__(self, data_loader: DataLoader, country: str):
        self.data_loader = data_loader
        self.country = country
    
    def calc(self, dt: str):
        return f"[{self.country}] 衍生资产数据: {dt}"

2. 工厂函数生成可复用资产

编写工厂函数,传入country参数即可生成该国家的所有资产,资产名称带国家标识避免冲突:

def create_country_assets(country: str):
    # 生成第一个资产
    @asset(
        name=f"my_asset_{country}",
        partitions_def=daily_partitions_def,
        config_schema=CountryAssetConfig
    )
    def my_asset(context: AssetExecutionContext):
        # 从配置中获取country(或直接用工厂函数传入的country,按需选择)
        target_country = context.config["country"]
        data_loader = DataLoader(country=target_country)
        return RunMyAsset(data_loader, target_country).calc(context.partition_key)
    
    # 生成第二个资产
    @asset(
        name=f"my_asset2_{country}",
        partitions_def=daily_partitions_def,
        config_schema=CountryAssetConfig
    )
    def my_asset2(context: AssetExecutionContext):
        target_country = context.config["country"]
        data_loader = DataLoader(country=target_country)
        return RunMyAsset2(data_loader, target_country).calc(context.partition_key)
    
    return [my_asset, my_asset2]

# 批量生成多国家资产
target_countries = ["JP", "US", "CN"]
all_assets = []
for country in target_countries:
    all_assets.extend(create_country_assets(country))

3. 为每个国家配置专属调度

同样用工厂函数生成调度,指定对应国家的资产组,并自定义调度时间适配时区:

from dagster import ScheduleDefinition, define_asset_job

def create_country_schedule(country: str):
    # 选择当前国家的所有资产
    asset_selection = [f"my_asset_{country}", f"my_asset2_{country}"]
    # 创建专属作业
    country_job = define_asset_job(
        name=f"daily_job_{country}",
        selection=asset_selection
    )
    # 按国家设置时区对应的cron表达式(示例:JP用东京时间0点,US用纽约时间0点)
    cron_map = {
        "JP": "0 0 * * *",  # UTC+9
        "US": "0 17 * * *", # UTC-5(对应纽约时间0点)
        "CN": "0 16 * * *"  # UTC+8(对应北京时间0点)
    }
    # 生成调度,同时传入country配置
    return ScheduleDefinition(
        job=country_job,
        cron_schedule=cron_map.get(country, "0 0 * * *"),
        config={
            "ops": {
                f"my_asset_{country}": {"config": {"country": country}},
                f"my_asset2_{country}": {"config": {"country": country}}
            }
        }
    )

# 批量生成所有调度
all_schedules = [create_country_schedule(c) for c in target_countries]

4. 注册资产与调度

在包的__init__.py中注册所有资产和调度:

from dagster import Definitions
from .assets import all_assets, all_schedules

defs = Definitions(
    assets=all_assets,
    schedules=all_schedules
)

关键优势

  • 所有逻辑集中在单个文件,无需为每个国家创建独立文件,维护成本低。
  • 资产与调度通过工厂函数批量生成,修改业务逻辑只需改一次。
  • 支持为每个国家自定义调度时间,适配不同时区需求。

内容的提问来源于stack exchange,提问作者pyCthon

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 18:14:59