使用 Argo Workflows 在 Kubernetes 上编排 Kedro 流水线:容器化、DAG 生成与集群部署实战
【免费下载链接】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 流水线迁移到 Kubernetes 集群、以容器方式并行执行节点任务的数据工程师。文章以仓库文档 docs/deploy/supported-platforms/argo.md 为主线,深入讲解如何将 Kedro 项目容器化、借助
build_argo_spec.py脚本把 Kedro 管道 DAG 转换为 Argo Workflows 规范、通过 Kubernetes Secrets 管理云存储凭证并完成提交与清理,同时结合 Kedro 源码剖析node_dependencies、bootstrap_project、kedro run -n等底层机制,帮助你理解"每个节点一个容器、任务间依赖由 Argo DAG 托管"的完整工作原理。
!!! warning "文档时效性声明"
仓库中 [docs/deploy/supported-platforms/argo.md](https://link.gitcode.com/i/dc3cd353e99158b958d2a9610f71e3b8) 明确标注该页面为 outdated documentation,尚未针对近期 Kedro 版本做过验证。因此本文给出的脚本与模板均以仓库现存文档为准,在使用较新版本的 Kedro 时,建议按文中步骤自行验证,并在遇到问题时到 Kedro 官方渠道反馈。为什么选择 Argo Workflows
Argo Workflows 是 Argo 项目四大组件之一,是一个开源、容器原生的工作流引擎,用于在 Kubernetes 上编排并行作业。它之所以适合作为 Kedro 流水线的执行平台,主要基于以下三点:
- 云无关:只要目标环境是任意一个 Kubernetes 集群,Argo Workflows 即可运行,不受特定云厂商绑定;
- 并行能力强:能够把计算密集型节点任务拆分到 Kubernetes 上并行运行与编排;
- DAG 依赖管理:使用有向无环图(DAG)管理任务之间的依赖关系,天然契合 Kedro 流水线的节点依赖结构。
从 Kedro 侧看,Pipeline.node_dependencies 属性正是以dict[Node, set[Node]]的形式暴露了每个节点及其父节点(直接依赖)集合——Argo Workflows 模板中的dependencies字段与之是一一对应的。这正是"Kedro 管道 DAG 到 Argo Workflows DAG"转换能够程序化完成的核心前提。
前置条件
在开始之前,请确认以下条件全部满足:
- Argo Workflows 已安装到你的 Kubernetes 集群(参见 Argo 官方 Quickstart);
- Argo CLI 已安装到你的本机,用于提交工作流;
- Kedro 流水线中每个节点都必须设置
name属性,因为节点名会被用来构建 DAG; - 所有节点的输入/输出数据集都必须在
catalog.yml中配置(参考 data_catalog_yaml_examples.md),并且必须指向外部存储位置(例如 AWS S3)。由于每个节点运行在独立容器中,节点之间无法通过进程内共享数据,因此不能在 workflow 中使用MemoryDataset。
!!! note "关于 MemoryDataset 的限制"
在 [DataCatalog](https://link.gitcode.com/i/1c7ad1b763aa4ba176275dacf42eb9de) 的实现中,`MemoryDataset` 是保存在 Python 进程内存中的数据集。单机 `SequentialRunner` 场景下节点间可以靠它传递数据,但 Argo Workflows 中每个节点各自运行在一个独立容器,进程不共享内存,内存数据集在节点边界处必然丢失,因此必须把中间数据落盘到 S3 等共享外部存储。这也是文档明确要求所有数据集配置外部位置的原因。!!! note "每个节点一个容器"
Argo Workflows 会为流水线中的**每一个节点启动一个独立的容器**(下文模板中每个 DAG task 都复用名为 `kedro` 的容器模板,仅通过参数切换要执行的节点)。第一步:容器化 Kedro 项目
首先需要用任意你偏好的容器方案(例如 Docker)把 Kedro 项目打包成镜像,供 Argo Workflows 使用。文档以 Docker 工作流为例,并推荐使用Kedro-Docker插件来简化镜像构建流程,具体步骤以该插件 README 为准。
镜像在本机构建完成后,需要把镜像推送到容器注册表(Container Registry),使 Kubernetes 集群能够拉取该镜像。相关说明可参考 docs/deploy/single_machine.md#how-to-use-container-registry。
第二步:程序化生成 Argo Workflows 规范
2.1 生成脚本 build_argo_spec.py
把下面的 Python 脚本保存到项目根目录(<project_root>/build_argo_spec.py)。它通过 Kedro 的公共 API 读取项目元数据与注册流水线,将节点的依赖关系转换为 Argo DAG 所需的任务列表,再结合 Jinja2 模板渲染出 Argo Workflows 规范:
# <project_root>/build_argo_spec.py import re from pathlib import Path import click from jinja2 import Environment, FileSystemLoader from kedro.framework.project import pipelines from kedro.framework.startup import bootstrap_project TEMPLATE_FILE = "argo_spec.tmpl" SEARCH_PATH = Path("templates") @click.command() @click.argument("image", required=True) @click.option("-p", "--pipeline", "pipeline_name", default=None) @click.option("--env", "-e", type=str, default=None) def generate_argo_config(image, pipeline_name, env): loader = FileSystemLoader(searchpath=SEARCH_PATH) template_env = Environment(loader=loader, trim_blocks=True, lstrip_blocks=True) template = template_env.get_template(TEMPLATE_FILE) project_path = Path.cwd() metadata = bootstrap_project(project_path) package_name = metadata.package_name pipeline_name = pipeline_name or "__default__" pipeline = pipelines.get(pipeline_name) tasks = get_dependencies(pipeline.node_dependencies) output = template.render(image=image, package_name=package_name, tasks=tasks) (SEARCH_PATH / f"argo-{package_name}.yml").write_text(output) def get_dependencies(dependencies): deps_dict = [ { "node": node.name, "name": clean_name(node.name), "deps": [clean_name(val.name) for val in parent_nodes], } for node, parent_nodes in dependencies.items() ] return deps_dict def clean_name(name): return re.sub(r"[\W_]+", "-", name).strip("-") if __name__ == "__main__": generate_argo_config()脚本参数说明:
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
image | 位置参数 | 是 | 已推送到容器注册表的镜像名 |
-p, --pipeline | 选项 | 否 | 要为其生成 Argo 规范的流水线名称;不传时默认使用__default__ |
-e, --env | 选项 | 否 | Kedro 配置环境名,默认是local |
脚本的底层原理(结合源码):
- bootstrap_project 负责"项目模式"下的启动设置:解析项目元数据、把
src目录加入sys.path与PYTHONPATH,并调用configure_project(package_name)完成项目配置;其返回值ProjectMetadata中的package_name被用于生成输出文件名argo-<package_name>.yml; - 从 kedro.framework.project 导入的
pipelines是一个懒加载的字典类对象,首次访问(如pipelines.get(...))时才会导入pipeline_registry模块并调用register_pipelines()读取已注册的流水线; pipeline.node_dependencies(见 pipeline.py)返回"节点 → 其父节点集合"的映射:父节点由其输出数据集被当前节点消费而建立依赖,独立节点对应空集合,在 Argo 模板中表现为无dependencies字段的并行任务;clean_name用正则[\W_]+把节点名中的非单词字符统一替换为-,保证生成的 Argo 任务名符合 Kubernetes 资源命名规范,同时 DAG 内的依赖引用也使用清洗后的名字保持一致。
2.2 Argo 规范模板 argo_spec.tmpl
将下面的模板保存到<project_root>/templates/argo_spec.tmpl。它定义了两个模板:kedro(单个节点容器执行模板)与dag(把每个 Kedro 节点映射为 DAG 任务):
{# <project_root>/templates/argo_spec.tmpl #} apiVersion: argoproj.io/v1alpha1 kind: Workflow metadata: generateName: {{ package_name }}- spec: entrypoint: dag templates: - name: kedro metadata: labels: {# Add label to have an ability to remove Kedro Pods easily #} app: kedro-argo retryStrategy: limit: 1 inputs: parameters: - name: kedro_node container: imagePullPolicy: Always image: {{ image }} env: - name: AWS_ACCESS_KEY_ID valueFrom: secretKeyRef: {# Secrets name #} name: aws-secrets key: access_key_id - name: AWS_SECRET_ACCESS_KEY valueFrom: secretKeyRef: name: aws-secrets key: secret_access_key command: [kedro]{% raw %} args: ["run", "-n", "{{inputs.parameters.kedro_node}}"] {% endraw %} - name: dag dag: tasks: {% for task in tasks %} - name: {{ task.name }} template: kedro {% if task.deps %} dependencies: {% for dep in task.deps %} - {{ dep }} {% endfor %} {% endif %} arguments: parameters: - name: kedro_node value: {{ task.node }} {% endfor %}模板要点逐项解读:
generateName: {{ package_name }}-:工作流名称由 Kedro 包名派生,避免命名冲突;- 容器入口是
kedro run -n <node>:即 kedro run 的--nodes / -n选项,它会把本次运行限制为指定名称的节点集合。于是每个 Argo DAG 任务通过参数kedro_node传入一个 Kedro 节点名,容器内部只执行该节点对应的函数,多个容器并行时即等价于把整条流水线并行化; retryStrategy: limit: 1:为每个节点容器配置了一次重试;app: kedro-argo标签:用于后续一键删除所有 Kedro 生成的 Pod(见清理小节);- 环境变量通过
secretKeyRef从名为aws-secrets的 Kubernetes Secret 注入 AWS 访问凭证; {% raw %}/{% endraw %}包裹args中的 Argo 占位符,避免与 Jinja2 语法冲突;- DAG 定义方式:任务的
dependencies直接来源于get_dependencies()生成的deps列表——这正是文档强调的"Argo Workflows 以有向无环图(DAG)定义任务间依赖关系"的落地形式。
!!! note "数据存储与 AWS 凭证"
本教程以 AWS S3 作为数据集存储。工作流要能读写 S3,需要把 `AWS_ACCESS_KEY_ID` 与 `AWS_SECRET_ACCESS_KEY` 注入容器;文档建议将两个值存入 Kubernetes Secrets(示例见下一节)。如果你的数据集存储在其他云厂商,将模板中的环境变量替换为对应凭证即可。2.3 安装 Jinja2 并运行生成脚本
模板使用 Jinja2 模板语言编写,因此需要先安装 Jinja 包:
$ pip install Jinja2然后在项目目录中运行辅助脚本生成 Argo 规范(生成的规范会保存到<project_root>/templates/argo-<package_name>.yml):
$ cd <project_root> $ python build_argo_spec.py <project_image>第三步:提交 Argo Workflows 规范到 Kubernetes
3.1 部署 Kubernetes Secret
提交工作流前,需要先部署一个 Kubernetes Secret 存放 AWS 凭证。以下是示例 Secrets 规范:
# secret.yml apiVersion: v1 kind: Secret metadata: name: aws-secrets data: access_key_id: <AWS_ACCESS_KEY_ID value encoded with base64> secret_access_key: <AWS_SECRET_ACCESS_KEY value encoded with base64> type: Opaque可以使用下面的命令把 AWS 密钥编码为 base64:
$ echo -n <original_key> | base64然后部署 Secret 到default命名空间并确认创建成功:
$ kubectl create -f secret.yml $ kubectl get secrets aws-secrets3.2 提交工作流
Secret 就绪后,即可提交 Argo 工作流:
$ cd <project_root> $ argo submit --watch templates/argo-<package_name>.yml!!! note "命名空间一致性"
提交 Argo Workflows 时,**必须保证与 Kubernetes Secrets 处于同一个命名空间**,否则容器无法通过 `secretKeyRef` 读取到凭证。更多用法参见 Argo CLI 帮助。3.3 清理集群资源
运行结束后,可以用以下命令清理集群资源(先删除带app=kedro-argo标签的 Pod,再删除 Secret):
$ kubectl delete pods --selector app=kedro-argo $ kubectl delete -f secret.yml备选方案:Kedro-Argo 插件
除手工编写脚本外,还可以使用社区提供的Kedro-Argo插件把 Kedro 项目转换为 Argo Workflows。
!!! warning "插件支持声明"
该插件**并非由 Kedro 官方团队维护**,仓库文档明确表示无法保证其可用性。若在正式生产环境中使用,请先充分评估并自行验证其行为。总结与实战要点
把 Kedro 流水线迁移到 Argo Workflows 的关键链路可以归纳为:
- 容器化:用 Docker(推荐 Kedro-Docker 插件)构建项目镜像并推送到容器注册表;
- 转换:借助
build_argo_spec.py+argo_spec.tmpl,从pipelines.get(name).node_dependencies提取节点依赖关系,渲染出argo-<package_name>.yml; - 凭证:用 base64 编码的 AWS 密钥创建
aws-secretsSecret; - 提交:
argo submit --watch templates/argo-<package_name>.yml,Argo 按 DAG 调度,每个节点在独立容器中以kedro run -n <node>执行; - 清理:按
app=kedro-argo标签删除 Pod,并删除 Secret。
实战中需要特别留意的三个约束:节点必须有name、数据集必须配置在catalog.yml且指向外部存储(禁止MemoryDataset)、工作流与 Secret 必须同命名空间。另外请牢记该部署文档本身标注为"outdated",在较新的 Kedro 版本上部署前务必重新验证脚本与模板的兼容性。更深入的流水线依赖机制可继续研读 Pipeline 源码 与 数据目录文档。
【免费下载链接】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),仅供参考