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

如何在脚本中访问已打包Kedro Pipeline的MemoryDataSet结果

在Kedro 0.18.9打包的自定义脚本中直接获取MemoryDataSet结果

核心思路

跳过main()的终止逻辑,手动加载打包项目的配置、Catalog,用Runner执行Pipeline后直接从MemoryDataSet提取DataFrame,完全避免文件读写或子进程调用。

步骤与代码示例

1. 打包前确保配置文件被包含

打包必须把conf目录纳入Python包,否则安装后无法加载Catalog配置:

  • 修改setup.py:
    setup(
        # 其他项目配置(名称、版本等)
        include_package_data=True,
        package_data={
            "your_package_name": ["conf/**/*"],
        },
    )
    
  • 或添加MANIFEST.in文件:
    include your_package_name/conf/**/*
    

2. 自定义脚本实现

from kedro.framework.project import configure_project
from kedro.config import ConfigLoader
from kedro.io import DataCatalog
from kedro.runner import SequentialRunner
from your_package_name.pipeline import create_pipeline  # 替换为你的Pipeline创建函数

# 配置Kedro项目(使用打包时的项目名称)
configure_project("your_package_name")

# 加载包内的Catalog配置
conf_loader = ConfigLoader(conf_source="your_package_name.conf")
catalog_config = conf_loader.get("catalog*", "catalog*/**")

# 初始化DataCatalog
catalog = DataCatalog.from_config(catalog_config)

# 获取要执行的Pipeline
pipeline = create_pipeline()

# 用SequentialRunner执行Pipeline(按需可替换为ParallelRunner)
runner = SequentialRunner()
runner.run(pipeline, catalog)

# 直接从MemoryDataSet提取结果DataFrame
# 替换为你的MemoryDataSet在catalog.yml中的名称
result_df = catalog.get("output_memory_dataset").data

# 后续自定义处理逻辑
print(result_df.head())

关键注意事项

  • 确认catalog.yml中输出数据集的配置为MemoryDataSet:
    output_memory_dataset:
      type: kedro.extras.datasets.pandas.MemoryDataSet
    
  • 如果Pipeline需要参数,需额外加载参数配置并传入Runner:
    params_config = conf_loader.get("parameters*", "parameters*/**")
    runner.run(pipeline, catalog, params=params_config)
    

内容的提问来源于stack exchange,提问作者trevinator

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 06:20:14