如何在脚本中访问已打包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
相关产品推荐
相关产品推荐

