Dagster运行Job时Op传入str类型参数报错的解决求助
问题解决:Dagster中调用Op传入字面量字符串报错
错误原因
Dagster的Job在组合Op时,要求Op的输入必须是上游Op的输出或者通过配置注入,不能直接传硬编码的字面量字符串(比如"Delivery_Location")。因为Job是声明式的依赖图,所有输入需要被纳入依赖管理体系,字面量无法被追踪、序列化或在运行时动态配置。
解决方法
方法1:使用Op配置(推荐)
把custom_table_name放到Op的配置类中,通过运行时配置传入:
修改Op定义:
from dagster import op, Config, get_dagster_logger # 定义配置类 class SaveToDBConfig(Config): custom_table_name: str @op def save_df_to_db(db: resource.sql_resource, dataframe, config: SaveToDBConfig): log = get_dagster_logger() table_name = config.custom_table_name log.info(f"Connecting to db, target table: {table_name}") with db.get().begin() as con: dataframe.to_sql(table_name, con, if_exists='append', index=False) # 其他逻辑...
修改Job调用:
@job def start_day_job(): delivery_location = import_ops.get_location() import_ops.save_df_to_db( dataframe=delivery_location, config={"custom_table_name": "Delivery_Location"} ) import_schedule = ScheduleDefinition( job=start_day_job, cron_schedule="0 7 * * 2-6",execution_timezone="Europe/Zurich" )
方法2:用静态Op返回字符串
如果不想用配置,可以创建一个返回固定字符串的Op,将其输出作为输入传给目标Op:
新增静态Op:
from dagster import op @op def get_delivery_location_table_name(): return "Delivery_Location"
修改Job定义:
@job def start_day_job(): delivery_location = import_ops.get_location() table_name = get_delivery_location_table_name() import_ops.save_df_to_db(dataframe=delivery_location, custom_table_name=table_name) import_schedule = ScheduleDefinition( job=start_day_job, cron_schedule="0 7 * * 2-6",execution_timezone="Europe/Zurich" )
说明
两种方法都能解决问题:方法1更灵活,支持运行时修改表名;方法2适合固定表名的场景,无需额外配置。
内容的提问来源于stack exchange,提问作者Ravi Teja
相关产品推荐
相关产品推荐

