尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Kedro 管道切片(Pipeline Slicing)完全指南:从输入、节点到标签与缺失输出的运行方式

Kedro 管道切片(Pipeline Slicing)完全指南:从输入、节点到标签与缺失输出的运行方式 Kedro 管道切片Pipeline Slicing完全指南从输入、节点到标签与缺失输出的运行方式【免费下载链接】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导读在实际数据工程中常常只需运行整条管道的一个子集slice例如跳过已完成的预处理、只重跑下游建模节点。本文基于 Kedro 的官方文档与当前仓库源码系统讲解两种切片入口——通过 Kedro-Viz 可视化切片与通过 Kedro CLI / Python API 编程式切片覆盖按输入、按节点、按结束节点、按标签、按指定节点以及只重建缺失输出六种方式并深入Pipeline与 Runner 的底层实现帮助你精确控制每一次kedro run的执行范围。切片的两条主路径Kedro 官方将管道切片slice a pipeline定义为只运行一条管道中部分节点a subset, or a slice。实现这一目标主要有两种方式通过 Kedro-Viz 可视化切片在 Kedro-Viz 界面中可视化地选择节点并生成切片随后在 Kedro 项目中执行生成的运行命令。这种方式适合交互式探索与人工决策详细操作步骤请参阅 Kedro-Viz 的《Slice a Pipeline》文档。通过 Kedro CLI 编程式切片向kedro run命令传递参数完整命令参考见 commands_reference.md或直接在 Python 中使用Pipeline类的切片方法。本文重点阐述这一路径。下图展示了 Kedro-Viz 中可视化切片的交互流程下面先从 pipeline_introduction.md 中如何构建管道一节的示例管道出发它计算一组数字的方差def mean(xs, n): return sum(xs) / n def mean_sos(xs, n): return sum(x**2 for x in xs) / n def variance(m, m2): return m2 - m * m full_pipeline pipeline( [ node(len, xs, n), node(mean, [xs, n], m, namemean_node, tagsmean), node(mean_sos, [xs, n], m2, namemean_sos, tags[mean, variance]), node(variance, [m, m2], v, namevariance_node, tagsvariance), ] )该管道共 4 个节点len由xs计算n、mean_node计算m、mean_sos计算m2、variance_node由m、m2计算v。调用Pipeline.describe()可以得到拓扑排序后的执行顺序与自由输入/输出#### Pipeline execution order #### Name: None Inputs: xs len([xs]) - [n] mean_node mean_sos variance_node Outputs: v ##################################从源码看describe() 返回执行顺序 自由输入 自由输出的格式化字符串其中节点顺序由 nodes 属性 中的拓扑排序决定基于graphlib.TopologicalSorter见 pipeline.py 的grouped_nodes实现。后续所有切片示例都将以describe()的输出为切片结果的判定依据。按输入切片from_inputs与--from-inputs一种切片方式是提供一组已预先计算好的输入让管道从这些输入开始向下游运行。例如希望从输入m2开始向下游切片print(full_pipeline.from_inputs(m2).describe())Output:#### Pipeline execution order #### Name: None Inputs: m, m2 variance_node Outputs: v ##################################切片后的管道只包含variance_node一个节点因为它同时需要m和m2而二者都没有上游节点产生在切片内所以它们都成为该切片管道的自由输入。再例如从输入m和xs同时切片print(full_pipeline.from_inputs(m, xs).describe())Output:#### Pipeline execution order #### Name: None Inputs: xs len([xs]) - [n] mean_node mean_sos variance_node Outputs: v ##################################注意把m加入from_inputs列表并不能保证m不会被重新计算——只要列表中同时指定了xsm的生产链路上游输入mean_node就会因直接/传递依赖xs而被纳入切片m自然会被重新生成。从源码看from_inputs() 的语义是直接或传递依赖这些输入的所有节点它先通过_get_nodes_with_inputs_transcode_compatible找到直接消费这些输入的节点然后以这些节点的输出作为新的起始输入沿_nodes_by_input索引循环向下游扩散广度优先传播直到没有新的节点可加入为止。CLI 中对应的选项是--from-inputs接受逗号分隔的输入数据集名见 project.py 中run命令的参数定义kedro run --from-inputsm2与之互补的是--to-outputs对应Pipeline.to_outputs见 project.py用于指定到哪些输出为止的切片。按起始节点切片from_nodes与--from-nodes另一种方式是指定一个或多个节点作为新管道的起点切片将包含该节点及其全部下游节点。例如print(full_pipeline.from_nodes(mean_node).describe())Output:#### Pipeline execution order #### Name: None Inputs: m2, n, xs mean_node variance_node Outputs: v ##################################该命令将管道从mean_node切片并运行它到所有下游节点。因为mean_node的输入是xs、n均为外部输入variance_node需要m由mean_node产生与m2外部输入所以切片管道的自由输入是m2、n、xs输出为v。在终端中运行对应的 CLI 命令kedro run --from-nodesmean_node从源码看from_nodes() 的实现其实是两个已有操作的组合先only_nodes(*node_names)取出指定节点再对其全部输出调用from_inputs(...)补齐所有下游节点res self.only_nodes(*node_names) res self.from_inputs(*map(_strip_transcoding, res.all_outputs()))因此from_nodes天然继承了from_inputs的传递下游语义。按结束节点切片to_nodes与--to-nodes与from_nodes相对可以指定一个或多个节点作为管道的终点切片将从管道开头运行到该节点为止含该节点及其全部上游依赖。例如print(full_pipeline.to_nodes(mean_node).describe())Output:#### Pipeline execution order #### Name: None Inputs: xs len([xs]) - [n] mean_node Outputs: m ##################################切片结果从开头运行到mean_node结束因为mean_node依赖xs与n其中n由len节点产生因此len节点被包含进来而mean_sos、variance_node被排除。对应 CLI 命令kedro run --to-nodesmean_node从源码看to_nodes() 与from_nodes对称先only_nodes(*node_names)再对其全部输入调用to_outputs(...)补齐所有上游节点res self.only_nodes(*node_names) res self.to_outputs(*map(_strip_transcoding, res.all_inputs()))同时指定起点与终点你还可以同时指定起始节点与结束节点从而精确圈定切片中包含的节点集合kedro run --from-nodesA --to-nodesZ当需要指定多个节点时用逗号分隔即可kedro run --from-nodesA,D --to-nodesX,Y,Z注意多个节点名之间使用英文逗号CLI 内部通过split_node_names回调拆分见 project.py。--from-nodes与--to-nodes对应的 help 文本分别定义为A list of node names which should be used as a starting point与...as an end pointproject.py。按标签切片only_nodes_with_tags与--tags还可以根据节点上的标签tags来切片。示例管道中mean_node带标签meanmean_sos同时带标签mean和variancevariance_node带标签variance。官方文档给出的同时拥有标签mean和variance的节点即mean_sos示例为print(full_pipeline.only_nodes_with_tags(mean, variance).describe())Output:#### Pipeline execution order #### Inputs: n, xs mean_sos Outputs: m2 ##################################而拥有标签mean或标签variance的节点官方示例通过两个调用的管道相加为Pipeline的合并运算实现sliced_pipeline full_pipeline.only_nodes_with_tags( mean ) full_pipeline.only_nodes_with_tags(variance) print(sliced_pipeline.describe())Output:#### Pipeline execution order #### Inputs: n, xs mean mean_sos variance Outputs: v ##################################实现语义说明从当前仓库源码看only_nodes_with_tags() 的实际筛选逻辑是集合交集判空——只要节点标签集合与给定标签集合存在任一交集unique_tags node.tags即选中该节点其 docstring 也明确写作包含anyof the provided tags 的节点即本质是 OR 语义unique_tags set(tags) nodes [node for node in self._nodes if unique_tags node.tags] return Pipeline(nodes)因此在实际运行中only_nodes_with_tags(mean, variance)会选中所有带mean或variance任一标签的节点文档第一个示例的输出描述与当前源码行为存在出入。如果你需要严格的 AND 语义同时包含全部指定标签可以在多个only_nodes_with_tags结果之间求交集或自行按节点标签过滤。CLI 对应选项是--tags/-t同样接受逗号分隔的多个标签见 project.py 中TAG_ARG_HELP定义kedro run --tagsmean,variance只运行指定节点only_nodes与--nodes当你需要精确运行管道中的某几个节点不附带它们的上下游依赖时使用only_nodesprint(full_pipeline.only_nodes(mean_node, mean_sos).describe())Output:#### Pipeline execution order #### Name: None Inputs: n, xs mean_node mean_sos Outputs: m, m2 ##################################该方法创建的新管道仅包含调用时指定的节点。CLI 对应--nodes/-n见 project.py 中NODE_ARG_HELPRun only nodes with specified nameskedro run --nodesmean_node,mean_sos重要提示由于切片没有附带这些节点的上游依赖指定节点的所有输入必须已经存在——要么已在数据目录data catalog中注册并可用要么由其他同时被选中的节点产生。否则运行时会因找不到输入而报错Runner 在运行前会校验pipeline.inputs()是否都能在 catalog 中满足见 runner.py 中对未满足输入的ValueError检查。从源码看only_nodes() 还会对不存在的节点名抛出ValueError并且当传入的名字以命名空间形式存在时如pipeline.only_nodes(model_training)而实际节点名为company.model_training会给出Did you mean的友好提示。重建缺失输出run_only_missing与--only-missing-outputsKedro 还能根据已有输出自动生成切片管道。当某些节点输出已持久化到磁盘时跳过这些节点可以避免重复运行耗时步骤。继续以示例管道为例先用JSONDataset保存中间输出nfrom kedro_datasets.pandas import JSONDataset from kedro.io import DataCatalog, MemoryDataset n_json JSONDataset(filepath./data/07_model_output/len.json) io DataCatalog(dict(xsMemoryDataset([1, 2, 3]), nn_json))由于n尚未保存检查其是否存在返回Falseio.exists(n)Output:Out[15]: False运行完整管道后n被计算并保存到磁盘SequentialRunner().run(full_pipeline, io)Output:Out[16]: {v: 0.666666666666667}io.exists(n)Output:Out[17]: True此时调用只运行缺失输出的入口文档示例为Runner.run_only_missing在当前仓库实现中该能力以run(..., only_missing_outputsTrue)参数的形式提供见 runner.py 的签名与参数说明CLI 对应--only-missing-outputs标志见 project.py可以跳过已经产出结果的第一个节点len([xs]) - [n]SequentialRunner().run(full_pipeline, io, only_missing_outputsTrue)Output:Out[18]: {v: 0.666666666666667}可以看到输出v依然正确但原始管道的第一个节点len产生n并未被重新执行。验证后清理临时文件try: os.remove(./data/07_model_output/len.json) except FileNotFoundError: pass底层机制逆拓扑序 缺失输出判定从源码看该能力由 AbstractRunner.run() 在收到only_missing_outputsTrue时触发_filter_pipeline_for_missing_outputsrunner.py实现其核心思路是将节点按逆拓扑序排列list(pipeline.nodes); sorted_nodes.reverse()基于输出数据集与生产节点的映射构建节点→子节点关系图_build_node_children_map逐个节点依据以下规则判定是否必须运行_should_node_run见 runner.py无输出如副作用节点的节点总是运行存在缺失的持久化输出catalog 中已注册、非_EPHEMERAL临时数据集、且catalog.exists()为 False则运行其输出被将要运行的子节点需要且该数据集缺失或为临时数据集则运行用筛选后的节点集合调用pipeline.filter(node_names...)生成切片管道并输出日志跳过多少个已有输出的节点、实际运行多少个节点见_log_filtering_resultsrunner.py。对应的 CLI 命令为kedro run --only-missing-outputs运行时日志会显示诸如Skipping 1 nodes with existing outputs: len和Running 3 out of 4 nodes之类的信息方便你确认跳过行为是否符合预期。切片方式的对照总览下表汇总了六种切片方式在 Python API 与 CLI 层的对应关系CLI 选项均定义于 project.py 的run命令切片方式Pipeline 方法CLI 选项语义按输入切片from_inputs(*inputs)--from-inputs从指定输入起含直接与传递依赖下游按输出切片to_outputs(*outputs)--to-outputs到指定输出止含直接与传递依赖上游按起始节点from_nodes(*names)--from-nodes从指定节点起含其全部下游节点按结束节点to_nodes(*names)--to-nodes到指定节点止含其全部上游节点按标签only_nodes_with_tags(*tags)--tags/-t选中含任一指定标签的节点OR 语义指定节点only_nodes(*names)--nodes/-n仅运行指定节点不附带上下游重建缺失输出run(..., only_missing_outputsTrue)--only-missing-outputs自动跳过输出已持久化的节点此外Pipeline还提供了统一的多条件过滤方法 filter()可一次传入tags、from_nodes、to_nodes、node_names、from_inputs、to_outputs、node_namespaces等条件与链式调用多个切片方法不同filter取所有条件的交集并且过滤结果为空时会抛出ValueError。对于大型模块化管道only_nodes_with_namespacespipeline.py和 CLI 的--namespaces选项还可以按命名空间批量圈选节点。总结与最佳实践优先使用--from-nodes/--to-nodes组合来圈定从 A 到 Z的连续子管道这是最直观、也最不容易出错的切片方式--nodes只适合输入已齐备的场景它不自动补齐上游务必确认所选节点的输入在数据目录中已存在或由同批节点产出--from-inputs不保证上游不重算只要切片中包含了输入的上游生产链路节点该输入仍会被重新计算--only-missing-outputs是增量重跑利器适合长时管道的断点续跑与增量更新但只对持久化输出缺失的节点生效临时ephemeral数据集不参与跳过判定标签切片默认是 OR 语义需要同时具备多个标签时应在多个only_nodes_with_tags结果间取交集或结合Pipeline.filter的条件交集语义实现。上述所有方法都定义在 kedro/pipeline/pipeline.py 的Pipeline类中CLI 参数解析与run命令实现位于 kedro/framework/cli/project.py底层运行调度含缺失输出过滤、拓扑排序执行与失败续跑建议见 kedro/runner/runner.py。你可以在自己的 Kedro 项目中用kedro run --help查看完整选项并借助Pipeline.describe()快速验证切片结果是否符合预期。【免费下载链接】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),仅供参考
返回列表