Kedro 节点(Node)完全指南:从基础定义到预览函数的实战进阶 📅 发布时间:2026/9/15 19:16:42 👁 浏览次数: Kedro 节点Node完全指南从基础定义到预览函数的实战进阶【免费下载链接】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/kedroKedro 中的节点Node是构建数据管道的核心构建单元它把一个普通 Python 函数与一组输入/输出数据变量绑定在一起形成管道中可调度、可追踪的最小任务。本文以 docs/build/nodes.md 为骨架结合 kedro/pipeline/node.py 与 kedro/pipeline/preview_contract.py 的源码实现系统讲解节点的创建语法、输入输出约定、*args/**kwargs函数、标签过滤、直接运行、生成器函数以及实验性的预览函数Preview Function功能。读完本文你将能够熟练编写、组合、调试并可视化 Kedro 节点并理解其底层校验与执行机制。节点是什么节点是管道的积木代表一项具体的任务。管道Pipeline通过组合节点构建工作流其范围可以是从简单的机器学习工作流到端到端E2E的生产级工作流。节点的相关 API 文档参见 kedro.pipeline.node。在继续之前请先导入 Kedro 及相关标准库from kedro.pipeline import * from kedro.io import * from kedro.runner import * import pickle import os从源码结构看kedro/pipeline/init.py 向外界导出了node、Node、Pipeline、GroupedNodes等符号因此上述通配符导入即可直接使用Node类。如何创建节点创建一个节点需要指定三个核心要素函数、输入变量名和输出变量名。以两个数相加的函数为例def add(x, y): return x y该函数有两个输入x和y和一个输出两者之和。使用它创建节点adder_node Node(funcadd, inputs[a, b], outputssum) adder_node输出如下Out[1]: Node(add, [a, b], sum, None)你还可以为节点添加name标签该标签会出现在日志中用于描述节点承载的业务逻辑adder_node Node(funcadd, inputs[a, b], outputssum) print(str(adder_node)) adder_node Node(funcadd, inputs[a, b], outputssum, nameadding_a_and_b) print(str(adder_node))输出如下add([a,b]) - [sum] adding_a_and_b: add([a,b]) - [sum]对照 node.py 中Node.__init__的签名逐项拆解上面的节点定义add节点运行时执行的 Python 函数参数func[a, b]输入变量名列表参数inputs个数必须与函数形参数一致sum返回值变量名参数outputsadd的返回值会被绑定到该变量上name可选的节点标签用于在日志和可视化中描述节点tags可选的节点标签集合用于按标签筛选运行confirms可选的待确认数据集名列表运行时会调用对应数据集的confirm()方法namespace可选的节点命名空间preview_fn可选的预览函数属于实验性功能详见下文。node.py中__str__的实现node.py正是通过拼接name: func(inputs) - outputs生成上述字符串这也是日志中节点描述的来源。节点定义语法节点的输入输出语法描述了一种约定它使得不同 Python 函数可以被复用进节点并支撑管道内部的依赖解析。输入变量语法输入语法含义示例函数形参节点运行时函数的调用方式None无输入def f()f()a单个输入def f(arg1)f(a)[a, b]多个输入def f(arg1, arg2)f(a, b)[a, b, c]可变输入def f(arg1, *args)f(arg1, arg2, arg3)dict(arg1x, arg2y)关键字输入def f(arg1, arg2)f(arg1x, arg2y)输出变量语法输出语法含义示例 return 语句None无输出不返回a单个输出return a[a, b]列表输出return [a, b]dict(key1a, key2b)字典输出return dict(key1a, key2b)以上任意组合都是允许的但Node(f, None, None)这种形式既无输入也无输出不合法——源码中 node.py 会直接抛出ValueError错误信息为it must have some inputs or outputs。从实现角度看Node.run()node.py会根据输入声明的形态分发到不同执行路径字符串走_run_with_one_input位置传参列表走_run_with_list按声明顺序展开为位置参数字典走_run_with_dict将数据集名映射为函数的关键字参数名。输出侧则由_outputs_to_dictionary统一将函数返回值整理成{输出变量名: 值}的字典。若输入声明与函数签名不匹配_validate_inputs还会通过inspect.signature(...).bind(...)在创建节点时即抛错实现早失败、早发现。*args节点函数在实际项目中经常需要处理任意数量的输入例如合并多个 DataFrame 的函数。你可以在节点函数中使用*args形参同时在节点inputs中声明这些数据集的名称Kedro 会按声明顺序依次传入def combine_all(*args): # 合并任意数量的 DataFrame ...配合列表形式的inputs例如Node(combine_all, [df1, df2, df3], combined)即可让同一函数灵活接收不同数量的输入。**kwargs专属节点函数有时比如编写上报/报告类节点你需要在节点内知道收到的数据集名称但事先并不确定有哪些。这时可以定义只接收**kwargs的函数def reporting(**kwargs): result [] for name, data in kwargs.items(): res example_report(name, data) result.append(res) return combined_report(result)随后构造Node时将节点输入声明为字典from kedro.pipeline import Node uk_reporting_node Node( reporting, inputs{uk_input1: uk_input1, uk_input2: uk_input2, ...}, outputsuk, ) ge_reporting_node Node( reporting, inputs{ge_input1: ge_input1, ge_input2: ge_input2, ...}, outputsge, )这里的字典键是函数的关键字参数名值是对应的数据集名。为避免重复手写这种键值同名的映射可以定义一个可复用的辅助函数from kedro.pipeline import Node mapping lambda x: {k: k for k in x} uk_reporting_node Node( reporting, - inputs{uk_input1: uk_input1, uk_input2: uk_input2, ...}, inputsmapping([uk_input1, uk_input2, ...]), outputsuk, ) ge_reporting_node Node( reporting, - inputs{ge_input1: ge_input1, ge_input2: ge_input2, ...}, inputsmapping([ge_input1, ge_input2, ...]), outputsge, )如何为节点打标签标签Tag可以在不改动代码的前提下筛选出管道的一部分来运行。例如kedro run --tagsds会运行所有带有ds标签的节点。为节点打标签只需在构造时传入tags参数Node(funcadd, inputs[a, b], outputssum, nameadding_a_and_b, tagsnode_tag)此外你也可以为整个 Pipeline 打标签若Pipeline(...)定义中包含tags参数Kedro 会把对应标签附加到该管道内的每一个节点上。按标签运行管道kedro run --tagspipeline_tag这会运行pipeline_tag标签下的全部节点。注意节点名和标签名只能包含字母、数字、连字符、下划线和句点其他符号一律不允许。这条规则在源码中有硬性校验__init__中通过正则[\w\.-]$分别校验name与每个tagnode.py不合法会抛出ValueError。tags属性返回的是去重后的set[str]且Node.tag(...)方法会返回一个追加了标签的节点副本因此标签操作不会修改原节点。在 CLI 层面kedro run的--tags/-t选项定义于 project.py支持逗号分隔多个标签由split_string回调拆分运行时过滤出匹配标签的节点子管道。如何运行节点运行节点前必须先实例化其输入。上例中的节点期望两个输入adder_node.run(dict(a2, b3))输出如下Out[2]: {sum: 5}提示你也可以像调用普通 Python 函数一样调用节点adder_node(dict(a2, b3))。这会在后台调用adder_node.run(dict(a2, b3))。这一点在源码中同样有体现Node.__call__直接委托给runnode.py因此两种调用方式完全等价tests/pipeline/test_node.py 中的test_call用例也验证了这一行为。同时run返回的是一个以输出变量名为键的字典——这正是 Kedro 管道中数据集Dataset流转的基础形态。如何在节点中使用生成器函数警告本节示例使用了pandas-irisstarter而该 starter 在 Kedro 0.19.0 及之后的版本中已不可用。最新支持它的版本是 Kedro 0.18.14请安装该版本或更早版本pip install kedro0.18.14来完成本示例。在终端输入kedro -V可查看当前 Kedro 版本。生成器函数Generator由 PEP 255 引入是一类返回惰性迭代器的特殊函数常用于数据的惰性加载或惰性保存非常适合处理无法一次性放入内存的大数据集。在 Kedro 中生成器函数可以用于节点内高效地分块处理大型数据。搭建项目使用旧版pandas-irisstarter 搭建项目假设 Kedro 版本为 0.18.14kedro new --starterpandas-iris --checkout0.18.14用生成器加载数据要在 Kedro 节点中使用生成器加载数据需要在catalog.yml中为相关数据集添加chunksize参数 X_test: type: pandas.CSVDataset filepath: data/05_model_input/X_test.csv load_args: chunksize: 10借助pandas的内置支持chunksize参数即可让数据以生成器的方式被分块读取。用生成器保存数据要使用生成器惰性保存数据需要做三件事将make_prediction函数定义从return改为yield创建名为ChunkWiseCSVDataset的自定义数据集更新catalog.yml以使用新创建的ChunkWiseCSVDataset。将以下代码复制到nodes.py。核心改动是改用DecisionTreeClassifier模型在make_predictions中按分块进行预测import logging from typing import Any, Dict, Tuple, Iterator, Generator from sklearn.preprocessing import LabelEncoder from sklearn.tree import DecisionTreeClassifier from sklearn.metrics import accuracy_score import numpy as np import pandas as pd def split_data( data: pd.DataFrame, parameters: Dict[str, Any] ) - Tuple[pd.DataFrame, pd.DataFrame, pd.Series, pd.Series]: Splits data into features and target training and test sets. Args: data: Data containing features and target. parameters: Parameters defined in parameters.yml. Returns: Split data. data_train data.sample( fracparameters[train_fraction], random_stateparameters[random_state] ) data_test data.drop(data_train.index) X_train data_train.drop(columnsparameters[target_column]) X_test data_test.drop(columnsparameters[target_column]) y_train data_train[parameters[target_column]] y_test data_test[parameters[target_column]] label_encoder LabelEncoder() label_encoder.fit(pd.concat([y_train, y_test])) y_train label_encoder.transform(y_train) return X_train, X_test, y_train, y_test def make_predictions( X_train: pd.DataFrame, X_test: pd.DataFrame, y_train: pd.Series ) - Generator[pd.Series, None, None]: Use a DecisionTreeClassifier model to make prediction. model DecisionTreeClassifier() model.fit(X_train, y_train) for chunk in X_test: y_pred model.predict(chunk) y_pred pd.DataFrame(y_pred) yield y_pred def report_accuracy(y_pred: pd.Series, y_test: pd.Series): Calculates and logs the accuracy. Args: y_pred: Predicted target. y_test: True target. accuracy accuracy_score(y_test, y_pred) logger logging.getLogger(__name__) logger.info(Model has accuracy of %.3f on test data., accuracy)ChunkWiseCSVDataset是pandas.CSVDataset的一个变体主要改动在_save方法首块数据覆盖写入含表头后续块追加写入。你需要新建src/package_name/chunkwise.py并把该类放入其中。参考实现如下import pandas as pd from kedro.io.core import ( get_filepath_str, ) from kedro_datasets.pandas import CSVDataset class ChunkWiseCSVDataset(CSVDataset): ChunkWiseCSVDataset loads/saves data from/to a CSV file using an underlying filesystem. It uses pandas to handle the CSV file. _overwrite True def _save(self, data: pd.DataFrame) - None: save_path get_filepath_str(self._get_save_path(), self._protocol) # Save the header for the first batch if self._overwrite: data.to_csv(save_path, indexFalse, modew) self._overwrite False else: data.to_csv(save_path, indexFalse, headerFalse, modea)随后更新catalog.yml让y_pred使用这个新数据集 y_pred: type: package_name.chunkwise.ChunkWiseCSVDataset filepath: data/07_model_output/y_pred.csv完成以上改动后在终端运行kedro run你会在日志中看到y_pred被多次保存——这正是生成器分块处理并保存数据的过程... INFO Loading data from y_train (MemoryDataset)... data_catalog.py:475 INFO Running node: make_predictions: make_predictions([X_train,X_test,y_train]) - [y_pred] node.py:331 INFO Saving data to y_pred (ChunkWiseCSVDataset)... data_catalog.py:514 INFO Saving data to y_pred (ChunkWiseCSVDataset)... data_catalog.py:514 INFO Saving data to y_pred (ChunkWiseCSVDataset)... data_catalog.py:514 INFO Completed 2 out of 3 tasks sequential_runner.py:85 INFO Loading data from y_pred (ChunkWiseCSVDataset)... data_catalog.py:475 ... runner.py:105从源码看生成器输出之所以能被节点正确消费是因为_outputs_to_dictionarynode.py对生成器返回值做了特殊处理它通过more_itertools.spy窥探生成器的第一个产出再在后续迭代时逐块收取结果并组装成输出字典。这也解释了为什么日志中会出现连续多次Saving data to y_pred。如何为节点添加预览函数警告该功能目前处于实验阶段可能在未来的版本中更改或移除。实验性功能遵循 docs/about/experimental.md 中描述的流程。预览函数Preview Function允许你注入一个可调用对象用于调试和监控。它不必加载完整数据集而是返回轻量级的摘要、代码片段、图表或示意图。概览预览函数是一个返回预览负载Preview Payload的可调用对象。预览负载可以是摘要Text用于日志和代码片段图表Mermaid用于关系和工作流图片Image用于绘图或可视化输出自定义格式Custom配合你自己的渲染器使用。预览函数通过节点的preview_fn参数挂载并可通过node.preview()调用。从源码看preview_fn在Node.__init__中必须可调用否则抛出ValueError首次使用时会触发一次KedroExperimentalWarning提示node.py。node.preview()node.py会调用该函数并通过isinstance校验返回值必须是四种预览类型之一否则抛出ValueError。基本用法from kedro.pipeline import node, Pipeline from kedro.pipeline.preview_contract import MermaidPreview import pandas as pd def train_model(training_data: pd.DataFrame) - dict: return { accuracy: 0.95, loss: 0.05, model_path: models/model_v1.pkl } def preview_training_model() - MermaidPreview: return MermaidPreview( content flowchart TD A[Training Started] -- B[Load Dataset] B -- C[Training Samples: 10,000] B -- D[Validation Samples: 2,000] C -- E[Train Model] D -- E E -- F[Epochs: 10] F -- G[Status: Completed] , meta{ timestamp: 2024-01-15T10:30:00, framework: sklearn } ) pipeline Pipeline( [ node( functrain_model, inputstraining_data, outputsmodel_metrics, # injecting a node preview callable preview_fnpreview_training_model, nametrain_model_node, ) ] ) # Get the node training_node next(n for n in pipeline.nodes) # Generate preview preview training_node.preview() # Returns MermaidPreview object preview_dict preview.to_dict() # Serialise for APIs/frontends可用的预览类型导入你需要的预览类型from kedro.pipeline.preview_contract import ( MermaidPreview, ImagePreview, TextPreview, CustomPreview, )这些类型全部定义于 preview_contract.py均为frozenTrue的 dataclass共享meta可选元数据字段并提供to_dict()方法将负载序列化为 JSON 安全的字典asdict JSON 可序列化校验便于供 API 或前端消费。CustomPreview还会额外校验renderer_key为非空字符串、content必须是 JSON 安全的字典。Mermaid 预览适用于图表、流程图或流程可视化def preview_pipeline_flow() - MermaidPreview: return MermaidPreview( content graph LR A[Load Data] -- B[Clean Data] B -- C[Feature Engineering] C -- D[Train Model] D -- E[Evaluate] )你还可以在meta参数中提供配置对象自定义 Mermaid 图表在 Kedro-Viz 中的渲染方式从而控制布局、样式、文本换行等选项def generate_mermaid_preview() - MermaidPreview: Generate a Mermaid diagram with custom configuration. This example demonstrates how to customize both the Mermaid rendering configuration and the text styling for node labels. diagram graph TD A[Raw Data] --|Ingest| B(Typed Data) B -- C{Quality Check} C --|Pass| D[Clean Data] C --|Fail| E[Error Log] D -- F[Feature Engineering] F -- G[Model Training] G -- H[Predictions] style A fill:#e1f5ff style D fill:#c8e6c9 style E fill:#ffcdd2 style H fill:#fff9c4 # Customize Mermaid rendering configuration # NOTE: On Kedro-Viz, this configuration will be # merged with sensible defaults custom_config { securityLevel: strict, # Security level: strict, loose, antiscript flowchart: { wrappingWidth: 300, # Text wrapping threshold (default: 250) nodeSpacing: 60, # Horizontal space between nodes (default: 50) rankSpacing: 60, # Vertical space between levels (default: 50) curve: basis, # Edge curve style: basis, linear, step }, themeVariables: { fontSize: 16px, # Font size for labels (default: 14px) }, # CSS styling for text nodes textStyle: { padding: 6px, # Internal padding in nodes (default: 4px) lineHeight: 1.3, # Line height for wrapped text (default: 1.2) textAlign: center, # Text alignment (default: center) } } return MermaidPreview(contentdiagram, metacustom_config) node( funcprocess_data, inputsraw_data, outputsprocessed_data, preview_fngenerate_mermaid_preview, namedata_processing_node, )关于 Mermaid 配置项的完整列表可参考 Mermaid configuration schema documentation。图片预览适用于图表、绘图或可视化输出支持 URL 或 data URIdef preview_correlation_matrix() - ImagePreview: # Can return a URL return ImagePreview( contenthttps://example.com/correlation_matrix.png ) # Or a data URI for inline images # return ImagePreview( # contentdata:image/png;base64,iVBORw0KGgo... # )文本预览适用于文本摘要或日志def preview_processing_log() - TextPreview: return TextPreview( contentProcessed 1,000 records\nRemoved 50 duplicates\nFilled 23 missing values )你还可以在meta参数中指定语言让代码片段在 Kedro-Viz 中以带语法高亮的形式展示def generate_code_preview() - TextPreview: Generate a code preview with syntax highlighting. code def calculate_metrics(data): \\\Calculate key performance metrics.\\\ import pandas as pd metrics { mean: data.mean(), median: data.median(), std: data.std() } return pd.DataFrame(metrics) # Example usage result calculate_metrics(my_dataframe) print(result) return TextPreview(contentcode, meta{language: python}) node( funccalculate_metrics, inputsdata, outputsmetrics, preview_fngenerate_code_preview, namemetrics_calculation_node, )meta参数接受language键来指定语法高亮的语言。Kedro-Viz 支持python、javascript和yaml三种高亮。自定义预览对于特殊的渲染需求可以使用CustomPreviewdef preview_custom_visualization() - CustomPreview: return CustomPreview( renderer_keymy_custom_renderer, content{ type: network_graph, nodes: [...], edges: [...] } )renderer_key用于标识由哪个前端组件负责渲染该预览。为预览添加元数据所有预览类型都支持通过meta参数附加可选元数据。meta有两个用途通用元数据添加上下文信息如版本、时间戳或数据来源渲染配置控制预览在 Kedro-Viz 中的展示方式例如 Mermaid 图表布局、语法高亮。示例添加通用元数据def preview_processing_log() - TextPreview: return TextPreview( contentProcessed 1,000 records\nRemoved 50 duplicates\nFilled 23 missing values, meta{created_by: admin} )示例渲染配置Mermaid 预览自定义图表布局与样式文本预览开启代码语法高亮。在预览函数中使用数据上下文预览函数无法直接访问节点的输入或输出它们是独立定义的函数。如果需要基于真实数据生成预览可以在预览函数内部使用闭包或访问数据集。使用闭包捕获上下文def make_preview_fn(graph_diagram): Create a preview function with captured context. def preview_fn() - MermaidPreview: return MermaidPreview(contentgraph_diagram) return preview_fn # In your pipeline creation sample_diagram graph TD A[Raw Data] --|Ingest| B(Typed Data) B -- C{Quality Check} C --|Pass| D[Clean Data] C --|Fail| E[Error Log] D -- F[Feature Engineering] F -- G[Model Training] G -- H[Predictions] style A fill:#e1f5ff style D fill:#c8e6c9 style E fill:#ffcdd2 style H fill:#fff9c4 node( funcprocess_data, inputsdata, outputsresult, preview_fnmake_preview_fn(sample_diagram) )最佳实践保持预览轻量预览函数应返回摘要而非完整数据集。需要数据集预览时请使用数据集预览功能Dataset Preview而非节点预览让预览快速执行避免在预览函数中执行昂贵的计算选用恰当的类型根据数据形态选择最匹配的预览类型添加元数据包含时间戳、版本或数据来源等上下文信息处理错误必要时用 try-except 包裹预览逻辑测试预览函数确保它们总是返回合法的预览对象。节点在管道中的位置节点本身并不关心执行顺序——这是 Pipeline 的职责。Kedro 根据每个节点声明的输入输出自动解析依赖关系决定节点的执行次序并将一个节点的输出作为下游节点的输入。因此掌握节点定义语法尤其是输入输出的声明形态是理解管道依赖解析、切片Slice、命名空间Namespace以及模块化管道Modular Pipeline等一系列高级特性的前提。相关文档可进一步阅读管道对象Pipeline objects节点如何组合成管道、依赖解析与执行顺序切片管道基于标签等条件抽取管道子集运行命名空间通过namespace与点号约定管理大型管道模块化管道一个管道一个目录的工程组织方式如何创建自定义数据集生成器保存示例中ChunkWiseCSVDataset的完整背景。节点测试用例可参考 tests/pipeline/test_node.py其中覆盖了输入输出的各种声明形态、__call__与run的等价性、重复输入、无输入/无输出节点等边界场景是理解节点行为的绝佳补充材料。【免费下载链接】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),仅供参考