Kedro 渐进式改造 Jupyter Notebook 实战:从 pandas 脚本到 Data Catalog 与 OmegaConfigLoader
【免费下载链接】kedroKedro is a toolbox for production-ready data science. It uses software engineering best practices to help you create data engineering and data science pipelines that are reproducible, maintainable, and modular.项目地址: https://gitcode.com/GitHub_Trending/ke/kedro
本文是一份基于 Kedro 官方仓库notebook-example的渐进式改造指南:从一段完全不使用 Kedro 的 pandas/sklearn 探索性分析代码出发,逐步引入 Kedro 的 Data Catalog(数据目录)管理数据加载、用 YAML 配置文件消除“魔法数”、用OmegaConfigLoader统一加载配置,最后将脚本重构为符合 Kedro 节点(node)理念的纯函数集合。读完本文,你将掌握在不创建完整 Kedro 项目的前提下,把 Kedro 的数据管理与配置能力一点一点“移植”进 notebook 的完整路径,并理解每一步背后的源码原理,为将来把代码整体迁入正式 Kedro 项目铺平道路。
本文配套的完整示例代码位于仓库的docs/integrations-and-plugins/notebooks_and_ipython/notebook-example/目录,其中add_kedro_to_a_notebook.ipynb是与本文对应的 Jupyter Notebook(.md版由 jupytext 双向同步生成,两者内容等价)。
一、示例背景:Kedro spaceflights 教程
Kedro spaceflights 教程 以一组.py文件构成的 Kedro 项目形式讲解 Kedro 的基础知识,其故事背景如下:
现在是 2160 年,太空旅游业蓬勃发展。全球有数千家航天飞机公司将游客送上月球再返回。你拿到了每家航天飞机提供的设施清单、客户评价以及公司信息。
任务:构建一个模型,预测每次前往月球及返程航班的票价。
本文的 notebook 示例沿用同一份业务数据(companies、reviews、shuttles),但以 notebook 而非 Kedro 项目的形式运行。运行本文示例代码需要先完成 Kedro 的安装配置,参见安装指南。
原始示例代码:完全不使用 Kedro
改造前的完整示例代码如下。首先加载数据(需要 pandas 与 openpyxl 依赖):
import pandas as pd companies = pd.read_csv("data/companies.csv") reviews = pd.read_csv("data/reviews.csv") shuttles = pd.read_excel("data/shuttles.xlsx", engine="openpyxl")然后进行数据预处理(把字符串布尔值转成真布尔、清洗百分比和货币字段、合并三张表):
# Data processing companies["iata_approved"] = companies["iata_approved"] == "t" companies["company_rating"] = ( companies["company_rating"].str.replace("%", "").astype(float) ) shuttles["d_check_complete"] = shuttles["d_check_complete"] == "t" shuttles["moon_clearance_complete"] = shuttles["moon_clearance_complete"] == "t" shuttles["price"] = ( shuttles["price"].str.replace("$", "").str.replace(",", "").astype(float) ) rated_shuttles = shuttles.merge(reviews, left_on="id", right_on="shuttle_id") model_input_table = rated_shuttles.merge(companies, left_on="company_id", right_on="id") model_input_table = model_input_table.dropna() model_input_table.head()接着是模型训练(这里把特征列、test_size=0.3、random_state=3全部硬编码在代码里):
# Model training from sklearn.model_selection import train_test_split X = model_input_table[ [ "engines", "passenger_capacity", "crew", "d_check_complete", "moon_clearance_complete", "iata_approved", "company_rating", "review_scores_rating", ] ] y = model_input_table["price"] X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.3, random_state=3) from sklearn.linear_model import LinearRegression model = LinearRegression() model.fit(X_train, y_train) model.predict(X_test)以及模型评估:
# Model evaluation from sklearn.metrics import r2_score y_pred = model.predict(X_test) r2_score(y_test, y_pred)这段代码可以正常运行,但存在明显问题:数据读取逻辑与业务逻辑混在一起、路径和“魔法数”散落各处、所有中间变量依赖 notebook 的执行顺序。接下来我们分五步将其逐步“Kedro 化”。
二、第一步:用 Kedro Data Catalog 管理数据加载
即使暂时不打算使用完整的 Kedro 项目,也可以先把 Kedro 的数据处理能力引入现有 notebook 项目。
Kedro 的 Data Catalog 是项目所有可用数据源的注册表(registry),它提供了一个独立的位置来声明项目所用数据集的具体细节。Kedro 内置了针对不同文件类型和文件系统的数据集实现(dataset),因此你不需要自己编写任何读写数据的逻辑。
Kedro 提供的数据集种类丰富,包括 CSV、Excel、Parquet、Feather、HDF5、JSON、Pickle、SQL Tables、SQL Queries、Spark DataFrames 等,分别基于 pandas、spark、networkx、matplotlib、yaml 等 API 实现。它依赖fsspec从多种数据存储读取和保存数据,包括本地文件系统、网络文件系统、云对象存储和 Hadoop;同时你可以在 load/save 操作中传入额外参数,并使用版本化(versioning)和凭据(credentials)来控制数据访问。
编写 catalog.yml
要开始使用 Data Catalog,需要一份catalog.yml来定义函数中会用到的数据集。示例文件夹中已经提供了这样一份文件(见 catalog.yml):
companies: type: pandas.CSVDataset filepath: data/companies.csv reviews: type: pandas.CSVDataset filepath: data/reviews.csv shuttles: type: pandas.ExcelDataset filepath: data/shuttles.xlsx每个顶层键(如companies)是数据集名称,type指定数据集实现类(pandas.CSVDataset、pandas.ExcelDataset等),filepath指定文件路径。这些声明与示例数据文件一一对应:数据文件位于示例目录的 data/ 下(companies.csv、reviews.csv、shuttles.xlsx)。
在 notebook 中加载 Data Catalog
读取catalog.yml并实例化 Data Catalog 的代码如下:
# Using Kedro's DataCatalog from kedro.io import DataCatalog import yaml # load the configuration file with open("catalog.yml") as f: conf_catalog = yaml.safe_load(f) # Create the DataCatalog instance from the configuration catalog = DataCatalog.from_config(conf_catalog) # Load the datasets companies = catalog.load("companies") reviews = catalog.load("reviews") shuttles = catalog.load("shuttles")从源码看,DataCatalog.from_config是一个类方法工厂,接收“以数据集名称为键、以构造参数(type、filepath等)为值”的字典,返回一个已就绪的DataCatalog实例,见 kedro/io/data_catalog.py。其load(ds_name, version=None)方法(见 kedro/io/data_catalog.py)按名称取出数据集并执行加载。
经过这一步替换后,上文 spaceflights notebook 中后续的数据处理和模型评估代码可以原样继续运行——因为companies、reviews、shuttles这三个变量名保持一致,只是数据来源从“手写 pandas 读取”变成了“由 Data Catalog 声明并加载”。
三、第二步:用 YAML 配置文件消除“魔法数”
3.1 用配置文件管理“魔法数字”
在写探索性代码时,硬编码数值固然省事,但长期来看会让代码难以维护。上面模型评估示例中的sklearn.model_selection.train_test_split()调用,就向test_size和random_state传入了硬编码值:
X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.3, random_state=3)良好的软件工程实践建议把“魔法数字”提取为命名常量,可以定义在文件顶部、工具文件中,也可以使用 yaml 这类格式。示例文件夹中的 params.yml 正是为此准备的:
# params.yml model_options: test_size: 0.3 random_state: 3在 notebook 中引用该文件的值为:
import yaml with open("params.yml", encoding="utf-8") as yaml_file: params = yaml.safe_load(yaml_file)test_size = params["model_options"]["test_size"] random_state = params["model_options"]["random_state"]同时把特征列提取成独立变量,使模型代码更加清晰:
features = [ "engines", "passenger_capacity", "crew", "d_check_complete", "moon_clearance_complete", "iata_approved", "company_rating", "review_scores_rating", ]X = model_input_table[features] y = model_input_table["price"]from sklearn.model_selection import train_test_split X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=test_size, random_state=random_state )剩余的模型评估代码可以照常运行:
from sklearn.linear_model import LinearRegression model = LinearRegression() model.fit(X_train, y_train) model.predict(X_test) from sklearn.metrics import r2_score y_pred = model.predict(X_test) r2_score(y_test, y_pred)3.2 把“魔法值”全部纳入配置文件
如果将“魔法数字”的概念推广到一般性的“魔法值”,会发现features这个变量可能也会被别处复用。把它从代码中提取到名为parameters.yml的配置文件后,得到:
# parameters.yml model_options: test_size: 0.3 random_state: 3 features: - engines - passenger_capacity - crew - d_check_complete - moon_clearance_complete - iata_approved - company_rating - review_scores_rating示例文件夹中同样提供了这份 parameters.yml,notebook 中引用方式如下:
import yaml with open("parameters.yml", encoding="utf-8") as yaml_file: parameters = yaml.safe_load(yaml_file) test_size = parameters["model_options"]["test_size"] random_state = parameters["model_options"]["random_state"]X = model_input_table[parameters["model_options"]["features"]] y = model_input_table["price"]from sklearn.model_selection import train_test_split X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=test_size, random_state=random_state )同样,剩余的模型评估代码可以照常运行(LinearRegression训练、r2_score评估)。
四、第三步:用 Kedro 配置加载器接管 YAML 读取
上面两步中,我们一直用yaml.safe_load手工读取配置文件。Kedro 提供了配置加载器(configuration loader)来抽象“从 yaml 文件加载值”的过程。OmegaConfigLoader是 Kedro 默认的配置加载器实现,即使没有完整 Kedro 项目也可以直接使用,用它替代yaml.safe_load后,配置读取逻辑将更加统一、可扩展。
4.1 用 OmegaConfigLoader 加载“魔法值”
用 Kedro 的OmegaConfigLoader加载parameters.yml的代码如下:
from kedro.config import OmegaConfigLoader conf_loader = OmegaConfigLoader(conf_source=".")conf_params = conf_loader["parameters"] test_size = conf_params["model_options"]["test_size"] random_state = conf_params["model_options"]["random_state"] X = model_input_table[conf_params["model_options"]["features"]] y = model_input_table["price"]from sklearn.model_selection import train_test_split X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=test_size, random_state=random_state )剩余的模型评估代码照常运行。
从源码看,OmegaConfigLoader的类文档明确说明了它的工作方式:“递归扫描conf_source所包含目录中的配置文件(扩展名为yaml、yml或json),通过 OmegaConf 加载并合并,最后以配置字典形式返回”,见 kedro/config/omegaconf_config.py。当多个配置文件出现同名顶层键时,位于不同目录的配置以后处理的路径覆盖先前的键(环境覆盖 base);位于同一目录时则抛出ValueError。
conf_loader["parameters"]的取值操作对应OmegaConfigLoader.__getitem__,见 kedro/config/omegaconf_config.py:它先检查 key 是否在config_patterns中(不在则抛KeyError),再按模式匹配并加载、合并所有对应配置文件;若没有任何配置文件匹配,则抛出MissingConfigException。在本例中,conf_source="."表示以 notebook 所在目录为配置根目录,conf_loader["parameters"]会加载当前目录下的parameters.yml。
4.2 用 OmegaConfigLoader 加载 Data Catalog
前面我们用yaml.safe_load读取catalog.yml并传给DataCatalog类:
# Using Kedro's DataCatalog from kedro.io import DataCatalog import yaml # load the configuration file with open("catalog.yml") as f: conf_catalog = yaml.safe_load(f) # Create the DataCatalog instance from the configuration catalog = DataCatalog.from_config(conf_catalog) # Load the datasets ...同样,也可以改用 Kedro 的OmegaConfigLoader配置加载器来初始化 Data Catalog。加载catalog.yml的代码如下:
# Now we are using Kedro's ConfigLoader alongside the DataCatalog from kedro.config import OmegaConfigLoader from kedro.io import DataCatalog conf_loader = OmegaConfigLoader(conf_source=".") conf_catalog = conf_loader["catalog"] # Create the DataCatalog instance from the configuration catalog = DataCatalog.from_config(conf_catalog) # Load the datasets companies = catalog.load("companies") reviews = catalog.load("reviews") shuttles = catalog.load("shuttles")至此,数据管理(Data Catalog)与配置管理(OmegaConfigLoader)都由 Kedro 统一接管,yaml.safe_load被完全替代。
五、第四步:把代码重构为函数
走到这一步,notebook 代码已经引入了 Kedro 的数据管理和配置加载,变得更易于复用。如果最终目标是把代码迁出 notebook、放进完整的 Kedro 项目,还可以更进一步。
Kedro 项目中的代码运行在一条或多条 pipeline 中,pipeline 是一系列“节点”(node),节点包装离散的函数。因此我们把代码逐步改造成函数形态。
5.1 先尝试一个大函数
一种朴素方案是把所有逻辑塞进一个函数:
# Use Kedro for data management and configuration from kedro.config import OmegaConfigLoader from kedro.io import DataCatalog conf_loader = OmegaConfigLoader(conf_source=".") conf_catalog = conf_loader["catalog"] conf_params = conf_loader["parameters"] # Create the DataCatalog instance from the configuration catalog = DataCatalog.from_config(conf_catalog) # Load the datasets companies = catalog.load("companies") reviews = catalog.load("reviews") shuttles = catalog.load("shuttles") # Load the configuration data test_size = conf_params["model_options"]["test_size"] random_state = conf_params["model_options"]["random_state"]def big_function(): #################### # Data processing # #################### companies["iata_approved"] = companies["iata_approved"] == "t" companies["company_rating"] = ( companies["company_rating"].str.replace("%", "").astype(float) ) shuttles["d_check_complete"] = shuttles["d_check_complete"] == "t" shuttles["moon_clearance_complete"] = shuttles["moon_clearance_complete"] == "t" shuttles["price"] = ( shuttles["price"].str.replace("$", "").str.replace(",", "").astype(float) ) rated_shuttles = shuttles.merge(reviews, left_on="id", right_on="shuttle_id") model_input_table = rated_shuttles.merge( companies, left_on="company_id", right_on="id" ) model_input_table = model_input_table.dropna() model_input_table.head() X = model_input_table[conf_params["model_options"]["features"]] y = model_input_table["price"] ################################## # Model training and evaluation # ################################## from sklearn.model_selection import train_test_split X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=test_size, random_state=random_state ) from sklearn.linear_model import LinearRegression model = LinearRegression() model.fit(X_train, y_train) model.predict(X_test) from sklearn.metrics import r2_score y_pred = model.predict(X_test) print(r2_score(y_test, y_pred))# Call the one big function big_function()诚实地讲,这段代码的可维护性相比之前并没有质的提升——它仍然依赖函数外部的全局变量,职责混杂。
5.2 拆分成符合“节点”理念的纯函数
更好的做法是拆成一组小函数,对应 Kedro 所设想的“pipeline 由节点组成”的形态。Kedro 的节点应当表现一致、可重复、可预期:给定相同的输入,节点总是返回相同的输出——这正是“纯函数”(pure function)的定义。节点/纯函数应当是小型、单一职责的函数,只完成一件具体的事。
于是我们把big_function中的数据逻辑拆成一组各司其职的数据处理函数,再把模型训练与评估代码拆成三个独立的数据科学函数:
#################### # Data processing # #################### import pandas as pd def _is_true(x: pd.Series) -> pd.Series: return x == "t" def _parse_percentage(x: pd.Series) -> pd.Series: x = x.str.replace("%", "") x = x.astype(float) / 100 return x def _parse_money(x: pd.Series) -> pd.Series: x = x.str.replace("$", "").str.replace(",", "") x = x.astype(float) return x def preprocess_companies(companies: pd.DataFrame) -> pd.DataFrame: companies["iata_approved"] = _is_true(companies["iata_approved"]) companies["company_rating"] = _parse_percentage(companies["company_rating"]) return companies def preprocess_shuttles(shuttles: pd.DataFrame) -> pd.DataFrame: shuttles["d_check_complete"] = _is_true(shuttles["d_check_complete"]) shuttles["moon_clearance_complete"] = _is_true(shuttles["moon_clearance_complete"]) shuttles["price"] = _parse_money(shuttles["price"]) return shuttles def create_model_input_table( shuttles: pd.DataFrame, companies: pd.DataFrame, reviews: pd.DataFrame ) -> pd.DataFrame: rated_shuttles = shuttles.merge(reviews, left_on="id", right_on="shuttle_id") model_input_table = rated_shuttles.merge( companies, left_on="company_id", right_on="id" ) model_input_table = model_input_table.dropna() return model_input_table ################################## # Model training and evaluation # ################################## from typing import Dict, Tuple from sklearn.linear_model import LinearRegression from sklearn.metrics import r2_score from sklearn.model_selection import train_test_split def split_data(data: pd.DataFrame, parameters: Dict) -> Tuple: X = data[parameters["features"]] y = data["price"] X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=parameters["test_size"], random_state=parameters["random_state"] ) return X_train, X_test, y_train, y_test def train_model(X_train: pd.DataFrame, y_train: pd.Series) -> LinearRegression: regressor = LinearRegression() regressor.fit(X_train, y_train) return regressor def evaluate_model( regressor: LinearRegression, X_test: pd.DataFrame, y_test: pd.Series ): y_pred = regressor.predict(X_test) print(r2_score(y_test, y_pred))按 Kedro 节点之间的数据流顺序调用这些函数:
# Call data processing functions preprocessed_companies = preprocess_companies(companies) preprocessed_shuttles = preprocess_shuttles(shuttles) model_input_table = create_model_input_table( preprocessed_shuttles, preprocessed_companies, reviews ) # Call model evaluation functions X_train, X_test, y_train, y_test = split_data( model_input_table, conf_params["model_options"] ) regressor = train_model(X_train, y_train) evaluate_model(regressor, X_test, y_test)注意几个关键改进:split_data直接接收conf_params["model_options"]字典作为参数,函数内部用parameters["test_size"]、parameters["random_state"]、parameters["features"]取值,配置与代码彻底解耦;每个函数只依赖自己的入参并返回明确结果,不再依赖 notebook 的全局执行状态。这种“输入 → 函数 → 输出”的形态与 Kedro 节点在概念上完全一致,后续迁移到 Kedro 项目时只需把函数包进node()并组装成 pipeline。
5.3 完整汇总代码
最后,把全部内容汇总到一个 notebook cell 中以便对照参考——可以把它与本文开头最初的 notebook 代码进行比较:
# Kedro setup for data management and configuration from kedro.config import OmegaConfigLoader from kedro.io import DataCatalog conf_loader = OmegaConfigLoader(conf_source=".") conf_catalog = conf_loader["catalog"] conf_params = conf_loader["parameters"] # Create the DataCatalog instance from the configuration catalog = DataCatalog.from_config(conf_catalog) # Load the datasets companies = catalog.load("companies") reviews = catalog.load("reviews") shuttles = catalog.load("shuttles") # Load the configuration data test_size = conf_params["model_options"]["test_size"] random_state = conf_params["model_options"]["random_state"] #################### # Data processing # #################### import pandas as pd def _is_true(x: pd.Series) -> pd.Series: return x == "t" def _parse_percentage(x: pd.Series) -> pd.Series: x = x.str.replace("%", "") x = x.astype(float) / 100 return x def _parse_money(x: pd.Series) -> pd.Series: x = x.str.replace("$", "").str.replace(",", "") x = x.astype(float) return x def preprocess_companies(companies: pd.DataFrame) -> pd.DataFrame: companies["iata_approved"] = _is_true(companies["iata_approved"]) companies["company_rating"] = _parse_percentage(companies["company_rating"]) return companies def preprocess_shuttles(shuttles: pd.DataFrame) -> pd.DataFrame: shuttles["d_check_complete"] = _is_true(shuttles["d_check_complete"]) shuttles["moon_clearance_complete"] = _is_true(shuttles["moon_clearance_complete"]) shuttles["price"] = _parse_money(shuttles["price"]) return shuttles def create_model_input_table( shuttles: pd.DataFrame, companies: pd.DataFrame, reviews: pd.DataFrame ) -> pd.DataFrame: rated_shuttles = shuttles.merge(reviews, left_on="id", right_on="shuttle_id") model_input_table = rated_shuttles.merge( companies, left_on="company_id", right_on="id" ) model_input_table = model_input_table.dropna() return model_input_table ################################## # Model training and evaluation # ################################## from typing import Dict, Tuple from sklearn.linear_model import LinearRegression from sklearn.metrics import r2_score from sklearn.model_selection import train_test_split def split_data(data: pd.DataFrame, parameters: Dict) -> Tuple: X = data[parameters["features"]] y = data["price"] X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=parameters["test_size"], random_state=parameters["random_state"] ) return X_train, X_test, y_train, y_test def train_model(X_train: pd.DataFrame, y_train: pd.Series) -> LinearRegression: regressor = LinearRegression() regressor.fit(X_train, y_train) return regressor def evaluate_model( regressor: LinearRegression, X_test: pd.DataFrame, y_test: pd.Series ): y_pred = regressor.predict(X_test) print(r2_score(y_test, y_pred)) # Call data processing functions preprocessed_companies = preprocess_companies(companies) preprocessed_shuttles = preprocess_shuttles(shuttles) model_input_table = create_model_input_table( preprocessed_shuttles, preprocessed_companies, reviews ) # Call model evaluation functions X_train, X_test, y_train, y_test = split_data( model_input_table, conf_params["model_options"] ) regressor = train_model(X_train, y_train) evaluate_model(regressor, X_test, y_test)六、进度回顾与下一步方向
至此,这个 notebook 已经完成了“Kedro 化”:
- 数据管理:用 Data Catalog(
catalog.yml+DataCatalog.from_config)替代手写 pandas 读取; - 配置管理:用 YAML 配置文件(
parameters.yml)集中管理test_size、random_state、features等魔法值,并用OmegaConfigLoader统一加载; - 代码结构:把脚本式代码重构为一系列单一职责的纯函数,与 Kedro 节点、pipeline 的概念对齐。
这些改进让代码在未来更容易复用和迁移。如果目标是最终把代码移出 notebook 并在完整的 Kedro 项目中运行,可以继续阅读 Kedro spaceflights 教程 学习如何把函数包装成节点、组装 pipeline;关于配置加载器的更多细节(如多环境base/local合并规则、config_patterns自定义、变量插值等),可参考 OmegaConfigLoader API 文档 及其源码 kedro/config/omegaconf_config.py;Data Catalog 的完整 API(load/save、版本化、凭据、校验开关等)见 DataCatalog 文档 与实现 kedro/io/data_catalog.py。
七、示例文件结构与运行前提
本文示例的完整目录结构如下(相对仓库根目录):
docs/integrations-and-plugins/notebooks_and_ipython/notebook-example/ ├── add_kedro_to_a_notebook.ipynb # 与本文对应的 Notebook ├── add_kedro_to_a_notebook.md # jupytext 同步的 Markdown 版本 ├── catalog.yml # Data Catalog 配置(companies/reviews/shuttles) ├── params.yml # 仅含 test_size/random_state ├── parameters.yml # 含 test_size/random_state/features └── data/ ├── companies.csv ├── reviews.csv └── shuttles.xlsx运行前提:
- 已按安装指南配置好 Kedro(对应
kedro包及pandas、openpyxl等依赖); - 使用示例时请保留整个
notebook-example文件夹(或克隆整个仓库),因为 notebook 依赖该文件夹中的catalog.yml、parameters.yml与data/下的数据文件; - 在 notebook 中执行时,确保当前工作目录是
notebook-example文件夹(conf_source="."即指向该目录),否则catalog.yml、parameters.yml将无法被OmegaConfigLoader找到; - 如果你希望在本地复现 notebook 与 Markdown 的双向同步,可安装
jupytext后执行jupytext --set-formats md,ipynb add_kedro_to_a_notebook.md重新生成 notebook。
从源码结构看,OmegaConfigLoader与DataCatalog的上述用法均不依赖 Kedro 项目上下文(无需conf/base、settings.py等),这正是“在 notebook 中渐进式使用 Kedro 特性”能够成立的根本原因。
【免费下载链接】kedroKedro is a toolbox for production-ready data science. It uses software engineering best practices to help you create data engineering and data science pipelines that are reproducible, maintainable, and modular.项目地址: https://gitcode.com/GitHub_Trending/ke/kedro
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考