如何在Dagster中从Op向Asset传参并按调度多参数运行
如何在Dagster调度中用不同参数运行Asset
核心问题修正与实现步骤
你的代码目前存在两个关键问题:一是multiple_num资产没有接收外部参数的逻辑,二是RunRequest未传递参数给资产,同时cleanup_directory的定义和调用不匹配。以下是具体解决方法:
给资产添加参数支持
先通过Dagster的Config类定义参数结构,让multiple_num能接收外部传入的参数:from dagster import asset, Config class MultipleNumConfig(Config): value: str @asset def get_a(): return 2 # 替换为你的实际返回值 @asset def get_b(): return 3 # 替换为你的实际返回值 @asset(config_schema=MultipleNumConfig) def multiple_num(get_a, get_b, config: MultipleNumConfig): # 使用传入的参数进行逻辑处理 print(f"当前运行参数: {config.value}") return get_a * get_b # 替换为你的实际业务逻辑修改Op传递参数到RunRequest
在get_values的RunRequest中,通过run_config给资产注入参数,同时修正cleanup_directory的定义:from dagster import op, RunRequest, AssetKey def cleanup_directory(value: str) -> str: # 替换为你的实际清理逻辑 print(f"清理目录: {value}") return "success" @op def get_values(): values = ['a','s','e'] for value in values: yield RunRequest( run_key=f"multiple_num_{value}", # 用value作为唯一标识,避免重复运行 run_config={ "assets": { "multiple_num": { "config": { "value": value } } } }, asset_selection = [AssetKey("multiple_num")] ) cleanup_directory(value)配置作业与调度
将作业和调度关联,实现定时触发带参数的资产运行:from dagster import job, ScheduleDefinition @job def run_values(): get_values() # 示例:每分钟触发一次调度,可根据需求修改cron表达式 run_values_schedule = ScheduleDefinition( job=run_values, cron_schedule="* * * * *", )验证运行
启动Dagster UI后,调度触发的每次运行都会给multiple_num传入不同的value参数,可在运行日志中查看参数的使用情况。
内容的提问来源于stack exchange,提问作者Hari
相关产品推荐
相关产品推荐

