Dagster使用自定义Config传参后下游资产配置错误求助
解决Dagster管道中Config参数传递的配置错误
问题描述
搭建Dagster管道时,希望通过Config让用户在运行时自定义参数(实际场景为SQL查询的起止日期)。但在将使用该参数的filter_data资产结果传递到下游filter_again资产时,触发以下配置错误:
dagster._core.errors.DagsterInvalidConfigError: Error in config for op Error 1: Missing required config entry "config" at the root. Sample config for missing entry: {'config': {'fruit_select': '...'}}
在Dagster UI中输入参数(如"Bananas")后,filter_data可正常使用参数,但下游资产调用时触发错误。需要实现参数定义与资产数据流转的需求。
复现代码
import pandas as pd import random from datetime import datetime, timedelta from dagster import asset, MaterializeResult, MetadataValue, Config, materialize @asset def generate_dataset(): # Function to generate random dates def random_dates(start_date, end_date, n=10): date_range = end_date - start_date random_dates = [start_date + timedelta(days=random.randint(0, date_range.days)) for _ in range(n)] return random_dates # Set seed for reproducibility random.seed(42) # Define the number of rows num_rows = 100 # Generate random data fruits = ['Apple', 'Banana', 'Orange', 'Grapes', 'Kiwi'] fruit_column = [random.choice(fruits) for _ in range(num_rows)] units_column = [random.randint(1, 10) for _ in range(num_rows)] start_date = datetime(2022, 1, 1) end_date = datetime(2022, 12, 31) date_column = random_dates(start_date, end_date, num_rows) # Create a DataFrame df = pd.DataFrame({ 'fruit': fruit_column, 'units': units_column, 'date': date_column }) # Display the DataFrame print(df) return df class fruit_config(Config): fruit_select: str @asset(deps=[generate_dataset]) def filter_data(config: fruit_config): df = generate_dataset() df2 = df[df['fruit'] == config.fruit_select] print(df2) return df2 @asset(deps=[filter_data]) def filter_again(): df2 = filter_data() df3 = df2[df2['units'] > 5] print(df3) return df3
问题原因
- 错误调用上游资产:
filter_again中直接调用filter_data(),但filter_data需要传入fruit_config参数,运行时未提供导致配置缺失错误。 - 不符合Dagster依赖逻辑:Dagster资产间的数据流转应通过参数注入实现,而非直接调用资产函数。
解决方案
修改资产定义,让下游资产通过参数接收上游输出,同时修正上游资产的依赖获取方式:
修改后的代码
import pandas as pd import random from datetime import datetime, timedelta from dagster import asset, MaterializeResult, MetadataValue, Config, materialize @asset def generate_dataset(): def random_dates(start_date, end_date, n=10): date_range = end_date - start_date random_dates = [start_date + timedelta(days=random.randint(0, date_range.days)) for _ in range(n)] return random_dates random.seed(42) num_rows = 100 fruits = ['Apple', 'Banana', 'Orange', 'Grapes', 'Kiwi'] fruit_column = [random.choice(fruits) for _ in range(num_rows)] units_column = [random.randint(1, 10) for _ in range(num_rows)] start_date = datetime(2022, 1, 1) end_date = datetime(2022, 12, 31) date_column = random_dates(start_date, end_date, num_rows) df = pd.DataFrame({ 'fruit': fruit_column, 'units': units_column, 'date': date_column }) print(df) return df class fruit_config(Config): fruit_select: str # 通过参数接收上游资产,移除手动deps配置 @asset def filter_data(generate_dataset: pd.DataFrame, config: fruit_config): df2 = generate_dataset[generate_dataset['fruit'] == config.fruit_select] print(df2) return df2 # 通过参数接收filter_data的输出,无需处理config @asset def filter_again(filter_data: pd.DataFrame): df3 = filter_data[filter_data['units'] > 5] print(df3) return df3
关键修改说明
- 上游资产注入:
filter_data不再直接调用generate_dataset(),而是通过函数参数接收该资产的输出,Dagster会自动识别并处理依赖关系。 - 下游接收结果:
filter_again通过参数获取filter_data的输出,无需自行调用函数,避免重复要求config参数的问题。 - 自动依赖识别:移除冗余的
deps参数,Dagster会根据函数参数自动构建资产依赖链。
修改后,在UI中输入fruit_select参数运行管道时,filter_data会正确使用参数过滤数据,下游filter_again直接处理上游传递的结果,不会触发配置错误。
内容的提问来源于stack exchange,提问作者jdsimkin04
相关产品推荐
相关产品推荐

