Kedro 模块化流水线实战:从 `kedro pipeline create` 到自定义模板的完整指南
Kedro 模块化流水线实战从kedro pipeline create到自定义模板的完整指南【免费下载链接】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 提供了一套完整的模块化流水线Modular Pipelines机制将逻辑上相互隔离、可复用的代码拆分为一个流水线一个文件夹的独立模块。本文基于仓库中的 docs/build/modular_pipelines.md 官方文档结合 kedro/framework/cli/pipeline.py 的 CLI 实现、kedro/templates/pipeline 默认模板与 kedro/pipeline/pipeline.py 的Pipeline源码系统讲解如何创建、填充、删除、复用模块化流水线以及如何使用自定义 Cookiecutter 模板与流水线专属依赖。读完本文你将掌握模块化流水线的完整创建流程、命名约定与自动发现机制并能按需定制自己的流水线模板。为什么需要模块化流水线在典型的 Kedro 项目中默认生成的单一主流水线会随着项目演进变得越来越复杂难以维护。官方文档给出的建议是将代码拆分为逻辑上相互隔离、可复用的不同流水线模块每个流水线最好组织在独立的文件夹中以便在项目内和项目之间复制与复用。核心理念可概括为八个字一个流水线一个文件夹one pipeline, one folder。Kedro 为这一理念提供了三方面的工具支持使用kedro pipeline create命令创建全新的空白流水线一套清晰的流水线创建与组织结构约定使用自定义新流水线模板的能力。使用kedro pipeline create创建新流水线创建模块化流水线的命令非常简单kedro pipeline create pipeline_name命令执行后会在项目中生成一套包含样板文件夹与文件的流水线骨架。为了方便你直接开始开发Kedro 会一并生成流水线专属的nodes.py、pipeline.py、参数文件parameters file以及对应的tests测试结构并补齐所需的__init__.py文件。生成后的目录结构如下├── conf │ └── base │ └── parameters_{{pipeline_name}}.yml -- Pipeline-specific parameters └── src ├── my_project │ ├── __init__.py │ └── pipelines │ ├── __init__.py │ └── {{pipeline_name}} -- This folder defines the modular pipeline │ ├── __init__.py -- So that Python treats this pipeline as a module │ ├── nodes.py -- To declare your nodes │ └── pipeline.py -- To structure the pipeline itself └── tests ├── __init__.py └── pipelines ├── __init__.py └── {{pipeline_name}} -- Pipeline-specific tests ├── __init__.py └── test_pipeline.py其中{{pipeline_name}}会被实际流水线名称替换。这套结构与仓库内置的 默认流水线模板 一一对应模板中的config/目录含parameters_{{ cookiecutter.pipeline_name }}.yml会被复制到项目conf目录tests/目录会被复制到项目tests/pipelines/下。流水线名称的合法性校验从源码看create_pipeline 命令实现 会对流水线名称进行严格的 Python 包名合法性校验_assert_pkg_name_ok见 kedro/framework/cli/pipeline.py#L48-L75不合法时直接抛出KedroCliError必须以字母或下划线开头长度至少为 2 个字符除首字符外只能包含字母、数字与下划线不能是 Python 关键字如class、def等。这是因为生成的流水线最终会作为 Python 模块被导入一个合法的包名是自动发现机制得以工作的前提。相关 CLI 选项命令还支持若干选项官方文档提示可通过kedro pipeline create --help查看完整列表。结合源码可以确认以下常用选项选项说明pipeline_name必填位置参数新流水线名称会经过包名校验-t, --template path指定自定义 Cookcookiecutter 模板路径会覆盖任何本地模板--skip-config跳过为新建流水线生成配置文件-e, --env env指定在哪个配置环境中创建流水线配置默认base其中--env的实际默认值来自settings.CONFIG_LOADER_ARGS.get(base_env, base)见 kedro/framework/cli/pipeline.py#L124-L125。若未跳过配置生成且指定环境不存在命令会报错退出。模板查找优先级create_pipeline 的模板解析逻辑 遵循明确的优先级命令行--template参数优先级最高click 会自动校验路径存在项目级模板project_root/templates/pipeline未传--template时优先使用全局默认模板kedro/templates/pipeline以上都不存在时回退到 Kedro 自带的默认模板。对应测试 tests/framework/cli/pipeline/test_pipeline.py#L86-L115 验证了项目级模板被采用而 test_create_pipeline_template_command_line_override 则验证了命令行-t参数优先于项目本地模板。命令行输出中会打印Using pipeline template at: path以提示实际使用的模板位置。填充流水线nodes.py与pipeline.py创建完成后pipeline.py中会有一段模板代码需要你填入实际的流水线逻辑# src/my_project/pipelines/{{pipeline_name}}/pipeline.py from kedro.pipeline import Pipeline def create_pipeline(**kwargs) - Pipeline: return Pipeline([])这里创建了一个返回Pipeline类实例的create_pipeline()函数。必须保持函数名为create_pipeline()因为 Kedro 依靠该函数名自动发现流水线否则就需要在流水线注册表中手动注册。自动发现的具体过程见 docs/build/pipeline_registry.mdfind_pipelines()遍历src/package_name/pipelines/目录对每个子目录依次执行——导入package_name.pipelines.pipeline_name模块、调用其中的create_pipeline()函数、校验返回对象是Pipeline实例。默认情况下任何一步失败只会告警并跳过该流水线find_pipelines(raise_errorsFalse)方便开发期部分流水线处于未完成状态也能运行项目生产环境可传find_pipelines(raise_errorsTrue)让错误在第一时间暴露。在填充pipeline.py之前官方建议把所有节点函数存放在nodes.py中。沿用文档示例将mean()、mean_sos()、variance()三个函数写入nodes.py# src/my_project/pipelines/{{pipeline_name}}/nodes.py 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随后在pipeline.py中用Node把这些函数组装成流水线# src/my_project/pipelines/{{pipeline_name}}/pipelines.py from kedro.pipeline import Pipeline, Node from .nodes import mean, mean_sos, variance # Import node functions from nodes.py located in the same folder def create_pipeline(**kwargs) - Pipeline: return Pipeline( [ Node(len, xs, n), Node(mean, [xs, n], m, namemean_node, tagstag1), Node(mean_sos, [xs, n], m2, namemean_sos, tags[tag1, tag2]), Node(variance, [m, m2], v, namevariance_node), ], # A list of nodes and pipelines combined into a new pipeline tagstag3, # Optional, each pipeline node will be tagged namespace, # Optional inputs{}, # Optional outputs{}, # Optional parameters{}, # Optional )Node的三个核心位置参数分别是节点函数、输入数据集列表、输出数据集列表name与tags则为可选的元数据tags支持单个字符串或字符串列表两种写法。Pipeline构造参数详解上述create_pipeline()展示了Pipeline构造函数的若干可选参数。对照源码 kedro/pipeline/pipeline.py#L133-L192 中的__init__签名可以确认其完整语义参数类型作用nodesIterable[Node \| Pipeline] \| Pipeline构成流水线的节点列表如果列表中包含Pipeline对象这些流水线会被展开其全部节点并入新流水线inputsstr \| set[str] \| dict[str, str]暴露为上游连接点的输入。传str/set时保持原输入名传dict时可将现有输入名映射为新名称。只允许引用流水线的自由输入outputsstr \| set[str] \| dict[str, str]暴露给下游连接点的输出。dict支持名称重映射既可以引用自由输出也可以暴露中间结果parametersstr \| set[str] \| dict[str, str]需要做命名空间的参数名集合可省略params:前缀tagsstr \| Iterable[str]应用到流水线全部节点的标签集合namespacestr给所有数据集名称添加的前缀inputs/outputs显式指出的名称与参数引用params:和parameters除外prefix_datasets_with_namespacebool默认True是否给节点的输入、输出与参数加命名空间前缀当仅用命名空间做部署分组而不想重命名数据集时可设为False其中tags既可以在节点级设置也可以在流水线级统一设置——流水线级的标签会应用到其中所有节点。namespace、inputs、outputs、parameters则服务于流水线的复用具体用法详见 Reuse pipelines with namespaces。删除已有流水线如果需要删除某个已存在的流水线使用kedro pipeline delete pipeline_name从 delete_pipeline 实现 可以看到该命令会同时清理三类产物源码目录src/package/pipelines/pipeline_name/测试目录tests/pipelines/pipeline_name/配置目录下匹配的parameters_name.yml与catalog_name.yml兼容旧的嵌套目录结构parameters/name.yml与catalog/name.yml。默认情况下命令会列出待删除路径并交互式确认可用-y/--yes跳过确认非交互执行若三者均不存在会报Pipeline name not found。删除完成后命令还会提示如果你曾在pipeline_registry.py的register_pipelines()中注册过该流水线需要手动移除相关引用。使用自定义流水线模板如果你希望用自定义的 Cookiecutter 模板生成流水线将模板保存在project_root/templates/pipeline即可——kedro pipeline create会默认采用项目内的这个自定义模板。也可以通过--template标志显式指定模板路径kedro pipeline create pipeline_name --template path_to_template如前文模板查找优先级所述--template指定的模板优先于项目本地模板。Kedro 一个项目只支持一个默认流水线模板如需多个模板可存放在独立文件夹中用--template逐一指定。编写自定义模板的约束官方文档对自定义模板提出了明确要求模板必须渲染为合法、可导入的 Python 模块顶层包含返回Pipeline对象的create_pipeline函数必须包含config与tests子目录Kedro 在创建流水线时会把它们复制到项目的config与tests目录config与tests目录必须与默认模板保持相同的布局不可定制——不过参数文件的内容和实际测试文件内容可以自由修改除此之外文件与文件夹的名称、结构均可按需定制。仓库内置的 默认流水线模板 可作为编写起点它展示了最小可用结构cookiecutter.json{{ cookiecutter.pipeline_name }}/内含__init__.py、nodes.py、pipeline.py、config/parameters_{{ cookiecutter.pipeline_name }}.yml、tests/__init__.py、tests/test_pipeline.py。流水线模板使用Cookiecutter渲染因此模板中必须包含cookiecutter.json。可参考默认模板的 cookiecutter.json{pipeline_name: default, kedro_version: {{ cookiecutter.kedro_version }}}实际创建时_create_pipeline会注入渲染上下文cookie_context {pipeline_name: name, kedro_version: kedro.__version__}以no_inputTrue方式调用 cookiecutter因此{{ cookiecutter.pipeline_name }}与{{ cookiecutter.kedro_version }}会被自动替换。默认模板生成的pipeline.py骨架如下见 模板 pipeline.py This is a boilerplate pipeline {{ cookiecutter.pipeline_name }} generated using Kedro {{ cookiecutter.kedro_version }} from kedro.pipeline import Node, Pipeline # noqa def create_pipeline(**kwargs) - Pipeline: return Pipeline([])与 Starter 模板嵌套使用如果你把自定义流水线模板内嵌到 Kedro Starter 模板中需要在Starter 的cookiecutter.json而不是流水线模板自身的cookiecutter.json中添加_copy_without_render: [templates]这样在从 Starter 创建新项目时Cookiecutter 不会尝试渲染templates目录从而保留流水线模板中的{{ cookiecutter.pipeline_name }}占位符供后续kedro pipeline create使用。提供流水线专属依赖一个流水线可以在本地声明独立的requirements.txt文件来记录其外部依赖。这些依赖需要手动用pip安装pip install -r requirements.txt这意味着模块化流水线不仅可以复用代码还可以携带自己的第三方依赖清单进一步增强跨项目复用的完整性。小结Kedro 的模块化流水线机制通过一个流水线一个文件夹的组织约定将可复用的业务逻辑从日益膨胀的主流水线中解放出来。本文覆盖了从kedro pipeline create的目录骨架与命名校验、nodes.py/pipeline.py的填充规范与create_pipeline()自动发现约定、Pipeline构造参数tags、namespace、inputs、outputs、parameters的底层语义到kedro pipeline delete、自定义 Cookcookiecutter 模板的优先级与约束、以及流水线专属requirements.txt依赖管理。结合 CLI 源码、默认模板 与 测试用例 中的实现证据你可以在此基础上进一步阅读 流水线注册与自动发现 与 命名空间复用构建出高内聚、低耦合、可跨项目复用的模块化数据流水线体系。【免费下载链接】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),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →