如何用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
相关产品推荐
相关产品推荐

