You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kedro与Streamlit集成:传递自定义DataCatalog运行流水线的方法

问题:在Streamlit中调用Kedro流水线并传递动态DataCatalog的正确方式

我正在集成Kedro(最新版)与Streamlit(最新版)构建数据处理应用,目前已实现从Streamlit调用Kedro流水线,但需要将Streamlit中加载的自定义数据以DataCatalog形式传入Kedro流水线处理。具体流程:

  1. 在Streamlit中通过文件上传器加载CSV/Excel文件并转为DataFrame
  2. 执行指定Kedro流水线,传递包含该DataFrame的自定义DataCatalog
  3. 由流水线处理数据并存储结果

但执行时遇到三个错误:

  1. TypeError: KedroSession.run() got an unexpected keyword argument 'extra_params'
  2. AttributeError: can't set attribute 'catalog'
  3. AttributeError: 'KedroContext' object has no attribute 'pipeline_registry'

请问如何正确实现这一需求?


解决方案

错误原因与修复说明

  1. extra_params参数报错:新版Kedro(>=0.18.0)已移除KedroSession.run()的extra_params参数,需通过params参数传递额外配置,或直接修改上下文的params属性。
  2. 无法设置catalog属性:Kedro的KedroContext对象的catalog是只读属性,不能直接赋值,需基于原有Catalog创建新的自定义Catalog并传入session.run()。
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 08:35:30