如何在Kedro中动态合并多个同结构数据集?
动态加载Kedro多模型输出并合并的最优方案
场景与需求
在多模型Kedro流水线中,每个模型生成结构一致的CSV输出,需要合并后做后处理得到最终结果。当前采用静态定义节点输入的方式,无法适配新增模型或选择性运行子集,希望通过parameters.yml中的valid_models列表,用Kedro原生功能动态加载指定模型的输出数据集,避免手动读取文件的非最优实现。
原生解决方案
1. 动态构建流水线依赖
在pipeline.py中读取配置参数,动态生成合并节点的输入数据集列表,无需硬编码模型名称:
from kedro.pipeline import Pipeline, node from kedro.config import ConfigLoader from kedro.framework.project import settings def create_pipeline(**kwargs): # 加载全局参数 conf_loader = ConfigLoader(settings.CONF_SOURCE) params = conf_loader.get("parameters*", "parameters*/**") valid_models = params["valid_models"] # 动态生成需要合并的数据集名称 input_datasets = [f"{model}.output" for model in valid_models] # 合并节点逻辑(替换为你的后处理代码) def combine(*model_outputs): import pandas as pd merged_df = pd.concat(model_outputs, ignore_index=True) # 这里添加你的后处理步骤,比如清洗、统计等 return merged_df # 创建动态输入的合并节点 combine_node = node( func=combine, inputs=input_datasets, outputs="final_output", name="combine_model_outputs" ) return Pipeline([combine_node])
2. 保持现有Catalog配置不变
原catalog.yml的模型输出配置无需修改,Kedro会自动根据动态生成的数据集名称加载对应的数据:
model1.output: type: pandas.CSVDataSet filepath: ${data_path}/output/model1/output.csv model2.output: type: pandas.CSVDataSet filepath: ${data_path}/output/model2/output.csv model3.output: type: pandas.CSVDataSet filepath: ${data_path}/output/model3/output.csv final_output: type: pandas.CSVDataSet filepath: ${data_path}/output/final_output.csv layer: output
3. 通过参数控制模型子集
在parameters.yml中维护valid_models列表,新增或移除模型只需修改此配置,无需调整流水线代码:
valid_models: - model1 - model2
方案优势
- 完全基于Kedro原生功能,遵循其数据流管理规范,自动处理数据缓存、版本化和依赖追踪
- 无需手动读取文件,避免路径硬编码和Kedro无法感知数据依赖的问题
- 扩展性强:新增模型仅需添加
catalog.yml配置和更新valid_models列表,代码无需改动 - 支持选择性运行:通过修改
valid_models即可指定要合并的模型子集
内容的提问来源于stack exchange,提问作者Nandha Kumar
相关产品推荐
相关产品推荐

