使用 Kubeflow Pipelines 编排 Apache Beam 预处理管道:完整实战指南

发布时间:2026/10/6 18:49:35

使用 Kubeflow Pipelines 编排 Apache Beam 预处理管道:完整实战指南 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本文以 Apache Beam 仓库内置的 ML 编排示例为核心系统讲解如何将 Beam 数据管道数据验证、预处理、模型验证与推理等封装为 Kubeflow PipelinesKFP组件并通过 DAG 串联成端到端机器学习工作流。读完本文你将掌握 KFP 组件定义YAML 接口 容器化实现、组件间输入输出传递、管道编译与提交运行三个关键环节并能直接复用仓库中可运行的完整示例。一、为什么需要用 KFP 编排 Beam 管道Apache Beam 提供统一的批处理和流处理编程模型在机器学习项目中可承担数据验证、数据预处理、模型验证、模型部署与推理等任务相关任务清单可参阅 学习文档。但一个完整的 ML 工作流不止于此还包含数据探索、特征工程、模型训练等环节并且对可复现性与可审计性有硬性要求——每一步的元数据metadata和产物artifact都必须被记录与追踪。Kubeflow 正是面向这类诉求的开源 MLOps 平台。它以 DAG有向无环图形式构建、部署和管理端到端 ML 管道统一负责各步骤的调度与执行步骤之间执行参数的传递步骤之间元数据与产物的传递。核心集成模型Beam 管道自身的 DAG 可以作为Kubeflow 管道 DAG 中的一个节点。也就是说你可以在 KFP 层面编排何时运行数据预处理、何时启动训练等流程级逻辑而每个节点内部由 Beam 负责分布式数据处理。这种DAG 套 DAG的分层结构让流程编排与数据计算各司其职。二、三步走将 Beam 管道嵌入 KFP 的整体流程仓库示例 ml-orchestration/kfp 呈现了完整的集成过程整体分为三步创建 KFP 组件为每个步骤如数据摄取、预处理、训练定义接口component.yaml并将实现容器化创建 KFP 管道连接各组件定义输入输出如何在组件之间交换编译并运行将管道编译为 JSON 文件通过 KFP client 提交到集群端点执行。示例仓库中的目录结构如下这也是推荐的项目组织方式kfp ├── pipeline.py # KFP 管道定义 编译 提交 ├── requirements.txt # kfp、google-cloud-aiplatform 等 ├── pipeline.json # 编译产物已生成的管道定义 └── components ├── ingestion # 数据摄取组件 │ ├── Dockerfile │ ├── component.yaml │ ├── requirements.txt │ └── src/ingest.py ├── preprocessing # Beam 预处理组件本文重点 │ ├── Dockerfile │ ├── component.yaml │ ├── requirements.txt │ └── src/preprocess.py └── train # 模型训练组件 ├── Dockerfile ├── component.yaml ├── requirements.txt └── src/train.py每个组件由两部分组成component.yaml定义组件的输入/输出参数接口Python 源文件包含实际业务逻辑对预处理组件而言即 Beam 管道代码。三、第一步创建 KFP 组件3.1 用 YAML 定义组件接口以预处理组件为例其接口定义位于 components/preprocessing/component.yaml核心结构如下name: preprocessing description: Component that mimicks scraping data from the web and outputs it to a jsonlines format file inputs: - name: ingested_dataset_path description: source uri of the data to scrape type: String - name: base_artifact_path description: base path to store data type: String - name: gcp_project_id description: ID for the google cloud project to deploy the pipeline to. type: String - name: region description: Region in which to deploy the Dataflow pipeline. type: String - name: dataflow_staging_root description: Path to staging directory for the dataflow runner. type: String - name: beam_runner description: Beam runner, DataflowRunner or DirectRunner. type: String outputs: - name: preprocessed_dataset_path description: target uri for the ingested dataset type: String implementation: container: image: your-docker-registry/preprocessing-image-name:latest command: [ python3, preprocess.py, --ingested-dataset-path, {inputValue: ingested_dataset_path}, --base-artifact-path, {inputValue: base_artifact_path}, --preprocessed-dataset-path, {outputPath: preprocessed_dataset_path}, --gcp-project-id, {inputValue: gcp_project_id}, --region, {inputValue: region}, --dataflow-staging-root, {inputValue: dataflow_staging_root}, --beam-runner, {inputValue: beam_runner}, ]需要重点理解的三个部分inputs/outputs声明组件对外暴露的参数。每个参数需给出name、description和type示例中均为String。上游组件的输出正是通过名称匹配注入到这里的输入implementation.container.image指定组件实现所对应的容器镜像需要替换为真实推送的镜像地址implementation.container.command容器启动命令。KFP 将输入输出参数以命令行参数的形式传给组件实现其中{inputValue: xxx}会在运行时被替换为对应输入的实际值{outputPath: xxx}则被替换为 KFP 为输出分配的文件路径。3.2 用 ArgumentParser 接收参数正因为 KFP 以命令行参数方式注入输入输出组件实现必须使用argparse.ArgumentParser解析。预处理组件 src/preprocess.py 的解析器与 YAML 中声明的参数一一对应def parse_args(): Parse preprocessing arguments. parser argparse.ArgumentParser() parser.add_argument( --ingested-dataset-path, typestr, helpPath to the ingested dataset, requiredTrue) parser.add_argument( --preprocessed-dataset-path, typestr, helpThe target directory for the ingested dataset., requiredTrue) parser.add_argument( --base-artifact-path, typestr, helpBase path to store pipeline artifacts., requiredTrue) parser.add_argument( --gcp-project-id, typestr, helpID for the google cloud project to deploy the pipeline to., requiredTrue) parser.add_argument( --region, typestr, helpRegion in which to deploy the pipeline., requiredTrue) parser.add_argument( --dataflow-staging-root, typestr, helpPath to staging directory for dataflow., requiredTrue) parser.add_argument( --beam-runner, typestr, helpBeam runner: DataflowRunner or DirectRunner., defaultDirectRunner) return parser.parse_args() if __name__ __main__: args parse_args() preprocess_dataset(**vars(args))注意--beam-runner带有默认值DirectRunner说明该组件既可以在本地直接运行也可以切换为DataflowRunner上云执行——runner 的选择完全由参数驱动。3.3 容器化组件实现每个组件目录下都有独立的 DockerfileFROM python:3.9-slim # (Optional) install extra dependencies # install pypi dependencies COPY requirements.txt / RUN python3 -m pip install --no-cache-dir -r requirements.txt # copy src files and set working directory COPY src /src WORKDIR /src构建镜像时需要先构建、推送镜像再把镜像地址写回component.yaml的image字段。预处理器依赖列表见 components/preprocessing/requirements.txt包括apache_beam[gcp]、requests、torch、torchvision、numpy、Pillow等——这决定了预处理容器内可用的 Beam 与图像处理能力。3.4 组件内部如何传递输出KFP v1 组件只能通过文件方式写出输出。以摄取组件 components/ingestion/src/ingest.py 为例它先把真实数据写入base_artifact_path下的时间戳命名文件再把该文件的路径写入 KFP 分配的输出文件# timestamp as unique id for the component execution timestamp int(time.time()) # create directory to store the actual data target_path f{base_artifact_path}/ingestion/ingested_dataset_{timestamp}.jsonl # if the target path is a google cloud storage path convert the path to the gcsfuse path target_path_gcsfuse target_path.replace(gs://, /gcs/) Path(target_path_gcsfuse).parent.mkdir(parentsTrue, exist_okTrue) with open(target_path_gcsfuse, w) as f: f.writelines([...]) # KFP v1 components can only write output to files. The output of this # component is written to ingested_dataset_path and contains the path # of the actual ingested data Path(ingested_dataset_path).parent.mkdir(parentsTrue, exist_okTrue) with open(ingested_dataset_path, w) as f: f.write(target_path)这里有两个值得借鉴的工程细节产物落盘与路径传递分离真实数据写入gs://对象存储输出文件里只保存路径字符串KFP 层面传递的是轻量引用gcsfuse 路径转换gs://xxx会被转换为/gcs/xxx这是为了让容器内进程能以本地文件系统方式访问 GCS是云上 KFP 的常见处理手法。预处理组件的输出同理preprocess_dataset中把最终 Avro 数据的目录写入preprocessed_dataset_path输出文件见 preprocess.pytimestamp time.time() target_path f{base_artifact_path}/preprocessing/preprocessed_dataset_{timestamp} # the directory where the output file is created may or may not exists # so we have to create it. Path(preprocessed_dataset_path).parent.mkdir(parentsTrue, exist_okTrue) with open(preprocessed_dataset_path, w) as f: f.write(target_path)四、第二步创建 KFP 管道连接组件管道定义位于 kfp/pipeline.py完整演示了加载组件 → 装饰器声明管道 → 连接组件与传递输出的流程。4.1 管道级参数解析pipeline.py本身也使用ArgumentParser接收管道级参数def parse_args(): Parse arguments. parser argparse.ArgumentParser() parser.add_argument( --gcp-project-id, typestr, helpID for the google cloud project to deploy the pipeline to., requiredTrue) parser.add_argument( --region, typestr, helpRegion in which to deploy the pipeline., requiredTrue) parser.add_argument( --pipeline-root, typestr, helpPath to artifact repository where Kubeflow Pipelines stores a pipelines artifacts., requiredTrue) parser.add_argument( --component-artifact-root, typestr, helpPath to artifact repository where Kubeflow Pipelines components can store artifacts., requiredTrue) parser.add_argument( --dataflow-staging-root, typestr, helpPath to staging directory for dataflow., requiredTrue) parser.add_argument( --beam-runner, typestr, helpBeam runner: DataflowRunner or DirectRunner., defaultDirectRunner) return parser.parse_args() # arguments are parsed as a global variable so # they can be used in the pipeline decorator below ARGS parse_args() PIPELINE_ROOT vars(ARGS)[pipeline_root]参数在模块加载时即被解析为全局变量以便在下方dsl.pipeline装饰器中使用——这是 KFP 管道定义的惯用做法。4.2 从 YAML 加载组件# load the kfp components from their yaml files DataIngestOp comp.load_component(components/ingestion/component.yaml) DataPreprocessingOp comp.load_component( components/preprocessing/component.yaml) TrainModelOp comp.load_component(components/train/component.yaml)kfp.components.load_component读取 YAML 定义并将其包装为可调用的操作Op工厂。4.3 用装饰器声明管道并串联组件dsl.pipeline( pipeline_rootPIPELINE_ROOT, namebeam-preprocessing-kfp-example, descriptionPipeline to show an apache beam preprocessing example in KFP) def pipeline( gcp_project_id: str, region: str, component_artifact_root: str, dataflow_staging_root: str, beam_runner: str): KFP pipeline definition. ingest_data_task DataIngestOp(base_artifact_pathcomponent_artifact_root) data_preprocessing_task DataPreprocessingOp( ingested_dataset_pathingest_data_task.outputs[ingested_dataset_path], base_artifact_pathcomponent_artifact_root, gcp_project_idgcp_project_id, regionregion, dataflow_staging_rootdataflow_staging_root, beam_runnerbeam_runner) train_model_task TrainModelOp( preprocessed_dataset_pathdata_preprocessing_task. outputs[preprocessed_dataset_path], base_artifact_pathcomponent_artifact_root)这里的连接方式即是 KFP 依赖关系的声明语法数据流依赖DataPreprocessingOp的ingested_dataset_path直接引用ingest_data_task.outputs[ingested_dataset_path]TrainModelOp的preprocessed_dataset_path引用预处理任务的同名输出——KFP 据此自动推断 DAG 中任务的执行顺序参数透传管道函数的形参如gcp_project_id、region会原样透传给下游组件执行顺序保证训练任务只有在预处理任务产出preprocessed_dataset_path之后才会启动。三个组件摄取 → 预处理 → 训练由此构成一条清晰的端到端 ML 流水线摄取组件生成 JSONL 数据集Beam 预处理组件将其清洗、转换并写出 Avro训练组件加载预处理结果并保存模型。五、第三步编译与提交运行5.1 编译为 JSON 管道定义if __name__ __main__: Compiler().compile(pipeline_funcpipeline, package_pathpipeline.json)KFP v2 的Compiler().compile()将装饰器声明的管道函数转换编译为 JSON 文件仓库中已包含编译产物 pipeline.json可作为参考。编译后的 JSON 是自包含的管道定义描述所有任务、容器镜像、参数绑定与依赖关系。5.2 提交到 KFP 端点run_arguments vars(ARGS) del run_arguments[pipeline_root] client kfp.Client() experiment client.create_experiment(KFP orchestration example) run_result client.run_pipeline( experiment_idexperiment.id, job_nameKFP orchestration job, pipeline_package_pathpipeline.json, paramsrun_arguments)提交环节的要点kfp.Client()默认连接本地 KFP 端点可通过参数指定远程端点create_experiment创建一个实验可类比为命名空间返回实验 IDrun_pipeline将编译产物pipeline.json提交执行params传入管道级参数此处显式删除了仅供编译使用的pipeline_root运行结果返回run_result可用于轮询运行状态或查询产物。运行pipeline.py时的完整命令形态参数与parse_args对应python3 pipeline.py \ --gcp-project-id YOUR_PROJECT_ID \ --region us-central1 \ --pipeline-root gs://your-bucket/kfp/artifacts \ --component-artifact-root gs://your-bucket/kfp/components \ --dataflow-staging-root gs://your-bucket/dataflow/staging \ --beam-runner DataflowRunner管道级依赖kfp1.8.13、google-cloud-aiplatform记录在 kfp/requirements.txt 中。六、Beam 管道在组件内部如何运行源码剖析预处理组件内的 Beam 管道是整个编排的核心其运行配置值得单独拆解见 preprocess.py# We use the save_main_session option because one or more DoFns in this # workflow rely on global context (e.g., a module imported at module level). pipeline_options PipelineOptions( runnerbeam_runner, projectgcp_project_id, job_namefpreprocessing-{int(time.time())}, temp_locationdataflow_staging_root, regionregion, requirements_file/requirements.txt, save_main_sessionTrue, ) with beam.Pipeline(optionspipeline_options) as pipeline: ( pipeline | Read input jsonlines file beam.io.ReadFromText(ingested_dataset_path) | Load json beam.Map(json.loads) | Filter licenses beam.Filter(valid_license) | Download image from URL beam.FlatMap(download_image_from_url) | Resize image beam.Map(resize_image, sizeIMAGE_SIZE) | Clean Text beam.Map(clean_text) | Serialize Example beam.Map(serialize_example) | Write to Avro files beam.io.WriteToAvro( file_path_prefixtarget_path, schema{ namespace: preprocessing.example, type: record, name: Sample, fields: [{ name: id, type: int }, { name: caption, type: string }, { name: image, type: bytes }] }, file_name_suffix.avro))几个关键设计点runner 由 KFP 参数驱动beam_runner从 KFP 输入注入既可以在本地用DirectRunner快速验证也可以切换到DataflowRunner在 Google Cloud Dataflow 上分布式执行配合project、region、temp_location等参数save_main_sessionTrue源码注释明确指出工作流中的若干 DoFn 依赖全局上下文如模块级导入因此必须开启主会话保存否则远程执行时函数定义与模块导入会丢失requirements_file/requirements.txtDataflow 启动 worker 时按此文件安装依赖与组件 Dockerfile 中COPY requirements.txt /的路径保持一致并行化的预处理逻辑ReadFromText读取 JSONL →json.loads解析 →Filter按图片许可证筛选 →FlatMap按 URL 下载图片失败时仅记日志并跳过→resize_image统一尺寸 →clean_text清洗字幕文本 →serialize_example序列化图片 →WriteToAvro输出 Avro 文件。整条链路以声明式 PTransform 串联天然具备分布式并行能力。此外仓库中另有同属 ML 编排主题的 tfx 目录展示 Apache Beam 与 TFX 的编排集成可作为对比学习材料该示例的总体说明见 ml-orchestration/README.md。七、实战要点与常见坑YAML 参数名与 argparse 参数名要保持一致KFP 通过{inputValue: name}注入命令行参数组件内的ArgumentParser必须用同名通常为连字符风格接收否则运行时会因缺少必需参数而失败输出必须写文件KFP v1 组件只能通过{outputPath: xxx}指向的文件传递输出且目录可能不存在务必先mkdir(parentsTrue, exist_okTrue)云存储路径要转换在容器内访问 GCS 时gs://前缀需替换为/gcs/gcsfuse 挂载点镜像地址要真实可拉取示例中的your-docker-registry/xxx:latest为占位符需替换为已推送的镜像参数作用域要分清管道级参数编译期与运行期都需要与组件级参数运行期注入分层传递运行提交时注意剔除仅编译需要的参数如示例中的pipeline_root。按照上述三步流程你便可以在 Kubeflow 中以标准方式编排 Beam 预处理管道实现数据摄取、预处理、模型训练的全流程自动化、可复现与可审计。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam AI/ML 管道实战指南用 MLTransform 统一预处理、用 RunInference 完成模型推理Apache Beam AI/ML 管道实战指南用 MLTransform 统一预处理、用 RunInference 完成模型推理 Apache Beam 的大数据批处理流处理数据工程使用 Kubeflow Pipelines 编排 Apache Beam 机器学习工作流从组件封装到流水线提交使用 Kubeflow Pipelines 编排 Apache Beam 机器学习工作流从组件封装到流水线提交 Apache Beam 提供了统一的批处理与流大数据批处理流处理数据工程Friend 开源项目中 OMI-Composio 插件用 Notion 沉淀记忆并导入 OMI 的完整实战指南Friend 开源项目中 OMI Composio 插件用 Notion 沉淀记忆并导入 OMI 的完整实战指南 Friendomi开源仓库中的 plug上一篇Apache Spark GraphX 图计算编程指南属性图、核心算子与图算法实战下一篇CANN ops-nn 算子库 aclnnGeluBackwardV2 接口详解GELU 反向梯度计算的两段式调用指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/10/6 18:49:35

选录音软件如何避开后台录音中断的问题

很多人在长时间会议、线上网课、持续访谈场景下,都会遇到录音中途莫名停止的情况。切出录音界面回复消息、刷内容,再切回来就发现录音进程已经终止,前面几十分钟的内容直接丢失,事后补录完全没有可能性。这类问题不是偶然故障&…

2026/10/6 19:49:39

PHP heredoc语法错误全解析:从报错定位到版本差异与避坑实践

做PHP开发这些年,一提到heredoc,我脑子里第一反应不是方便,而是那条让人头大的“Parse error: syntax error”。尤其项目里用到heredoc字符串、邮件模板、批量SQL拼接时,代码动不动就报语法错误,很多时候明明看着缩进都…

2026/10/6 19:49:39

MOSFET体二极管反向恢复:双脉冲测试与仿真验证实战

1. 为什么体二极管反向恢复值得单独拎出来讲 做电源的同行大概都有过这种经历:板子焊好,上电,波形看着挺正常,效率也凑合,结果一跑满载或者一上高温,上下管直通炸机。拆下来一测,死区时间明明留…

2026/10/6 19:49:39

AI Agent技能包实战:从npx到GKE的skills设计与部署

1. 从“skills”这个标题说起:它到底指什么 第一次看到“skills”这个标题,很多人会以为是某个泛泛而谈的能力清单,或者一份简历上的技能罗列。但结合热词里的 Google Cloud、Agent Skills、npx、GKE、claude agent skills、codex skills 这些…

2026/10/6 19:49:39

功率MOSFET雪崩效应与UIS测试:从原理到选型实战

1. 从一次炸管说起:为什么MOSFET的雪崩效应值得单独拎出来讲功率MOSFET在开关电源、电机驱动、逆变器这些场景里,绝大多数失效都不是因为导通损耗算错了,而是因为关断瞬间的电压尖峰把器件打进了雪崩击穿区。我见过太多硬件工程师在调试阶段反…

2026/10/5 6:32:56

Jev+Agent接管浏览器:browser-use实战与jev-ultrafast性能优化

1. 从“Jev”说起:为什么我要把Agent接进浏览器“Jev”这个词最近在圈子里出现的频率越来越高,很多人第一次听到会以为是某个新模型的名字,其实它更像是一种思路——把Jev模型的能力当作底座,通过Agent的方式去接管浏览器&#xf…

2026/10/6 4:01:51

多智能体集群实战:DeepAgents编排、MCP与A2A协议及Skills体系

1. 从"单兵作战"到"集群协同":多智能体编排到底在解决什么问题如果你最近在折腾 Agent 相关的东西,大概率会有一种感觉:单个 Agent 能做的事情,其实很快就摸到天花板了。你给它一个提示词,挂几个工…

2026/10/6 17:46:51

无源低通滤波器设计实战:从RC到LC,手把手教你避开那些坑

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/6 0:03:23

MR25H40CDF+STM32F031C6工业级高可靠数据存储方案

1. 项目概述:为什么在工业现场非得用 MR25H40CDF 配 STM32F031C6 做数据存储?在工厂产线的 PLC 控制柜里、在风电变流器的散热片背面、在矿井监测终端的金属外壳下,你经常能看到一块指甲盖大小的黑色芯片——它既不是 Flash,也不是…

2026/10/6 0:03:23

MRAM+STM32工业断电数据保全实战指南

1. 项目概述:为什么在工业现场非得用 MR25H40CDF 配 STM32F031C6 做数据存储?在工厂产线的PLC柜里、在野外无人值守的环境监测终端里、在高速运转的包装机控制板上,你经常能看到一块指甲盖大小的黑色芯片,旁边贴着“MR25H40CDF”丝…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
☎咨询二维码 ☎ ↑