Kedro与Streamlit集成:传递自定义DataCatalog运行流水线的方法
问题:在Streamlit中调用Kedro流水线并传递动态DataCatalog的正确方式
我正在集成Kedro(最新版)与Streamlit(最新版)构建数据处理应用,目前已实现从Streamlit调用Kedro流水线,但需要将Streamlit中加载的自定义数据以DataCatalog形式传入Kedro流水线处理。具体流程:
- 在Streamlit中通过文件上传器加载CSV/Excel文件并转为DataFrame
- 执行指定Kedro流水线,传递包含该DataFrame的自定义DataCatalog
- 由流水线处理数据并存储结果
但执行时遇到三个错误:
TypeError: KedroSession.run() got an unexpected keyword argument 'extra_params'AttributeError: can't set attribute 'catalog'AttributeError: 'KedroContext' object has no attribute 'pipeline_registry'
请问如何正确实现这一需求?
解决方案
错误原因与修复说明
extra_params参数报错:新版Kedro(>=0.18.0)已移除KedroSession.run()的extra_params参数,需通过params参数传递额外配置,或直接修改上下文的params属性。- 无法设置
catalog属性:Kedro的KedroContext对象的catalog是只读属性,不能直接赋值,需基于原有Catalog创建新的自定义Catalog并传入session.run()。 pipeline_registry不存在:新版Kedro将pipeline_registry替换为pipelines属性,直接通过context.pipelines访问已注册的流水线。
完整实现代码
import streamlit as st import pandas as pd from pathlib import Path from kedro.framework.session import KedroSession from kedro.framework.startup import bootstrap_project from kedro.io import DataCatalog, MemoryDataSet # 配置Kedro项目路径(替换为你的实际项目路径) PROJECT_PATH = Path(__file__).parent.parent # 初始化Kedro项目环境 bootstrap_project(PROJECT_PATH) # Streamlit UI部分 st.title("Kedro + Streamlit 动态数据处理") uploaded_file = st.file_uploader("上传CSV/Excel文件", type=["csv", "xlsx"]) if uploaded_file: # 加载上传文件为DataFrame if uploaded_file.name.endswith(".csv"): df_uploaded = pd.read_csv(uploaded_file) else: df_uploaded = pd.read_excel(uploaded_file) st.subheader("上传的数据预览") st.dataframe(df_uploaded) # 创建Kedro会话并执行流水线 with KedroSession.create(project_path=PROJECT_PATH) as session: # 加载Kedro上下文 context = session.load_context() # 1. 基于原有Catalog创建自定义Catalog # 保留原有配置中的数据集,添加动态上传的数据集 custom_catalog = DataCatalog.from_config(context.config_loader.get("catalog")) # 将上传的DataFrame作为MemoryDataSet添加到Catalog,命名需与流水线中的数据集名称对应 custom_catalog.add("user_uploaded_data", MemoryDataSet(data=df_uploaded)) # 2. 指定要运行的流水线(替换为你的流水线名称,比如__default__) target_pipeline = context.pipelines["__default__"] # 3. 运行流水线,传入自定义Catalog与参数(如有需要) session.run( pipeline_name="__default__", catalog=custom_catalog, # 如需传递额外参数,可合并原有参数后传入 # params={**context.params, "custom_param": "your_value"} ) st.success("Kedro流水线执行完成!")
关键注意事项
- 项目路径配置:确保
PROJECT_PATH指向你的Kedro项目根目录,bootstrap_project需在会话创建前执行,以初始化Kedro环境。 - 数据集命名:自定义Catalog中添加的数据集名称(如
user_uploaded_data)需与Kedro流水线中定义的输入数据集名称完全一致,否则流水线无法识别。 - MemoryDataSet使用:动态加载的DataFrame通过
MemoryDataSet存储,无需在Kedro的catalog.yml中预先配置,适合临时数据传递。 - 参数传递:如需传递额外参数,直接通过
session.run()的params参数传入,格式为字典,可与原有上下文参数合并。
内容的提问来源于stack exchange,提问作者Riuk2252
相关产品推荐
相关产品推荐

