ZenML 第一条 AI 流水线实战指南:从 AI Agent、经典机器学习到混合系统
【免费下载链接】zenmlZenML 🙏: One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml
导读
本文基于 ZenML 官方入门文档《Your First AI Pipeline》整理编写,面向首次接触 ZenML 的开发者,系统讲解如何用同一套 steps → pipeline → artifacts → stacks 模式快速搭建三类流水线:LLM 驱动的 AI Agent 流水线、基于 scikit-learn/TensorFlow/PyTorch 的经典机器学习流水线,以及"分类器 + Agent"的混合系统。读完本文,你将掌握每条路径对应的仓库示例、核心源码结构、部署为实时 HTTP 服务的完整命令,以及通过 ZenML 仪表盘观察运行血缘(lineage)与步骤元数据的方法。
为什么 ZenML 流水线对三种场景一视同仁
ZenML 流水线对经典机器学习、AI Agent和混合方法采用完全相同的抽象,区别只在于你的 step 内部做什么、以及如何编排它们。官方入门文档归纳了三个核心理由:
- 可复现、可移植(Reproducible & portable):同一份代码,通过切换 stack(本地、远程、云端)即可在不同执行环境中运行,无需改动业务逻辑;
- 一套方法同时适用于模型与 Agent(One approach for models and agents):step、pipeline、artifact 的概念对 sklearn、经典 ML 与 LLM 同样成立;
- 默认可观测(Observe by default):血缘关系与步骤元数据(延迟 latency、token 消耗、指标 metrics 等)自动被追踪,并可在仪表盘中可视化。
值得留意的是,官方还推荐使用 Claude Code、Codex、Copilot、Cursor 等 AI 编码工具配合 ZenML 提供的 Agent Skills(zenml-scoping负责把想法拆解为多流水线计划,zenml-pipeline-authoring负责实现 step 与 pipeline),安装方式见 LLM tooling。不过即便完全手写,三条路径的骨架也是一致的。
从仓库源码看,最小骨架可以在 examples/quickstart 中看到:一个@step装饰的普通函数 simple_step.py,被@pipeline装饰的 simple_pipeline.py 调用,返回的字符串自动成为被追踪的 artifact。这个模式就是下面所有示例的最小单元。
Path 1:构建 AI Agent
用大语言模型、提示词和工具构建能够推理、采取行动并与你的系统交互的自主智能体。
两种运行 Agent 的方式
官方文档特别澄清了一个关键取舍:本条路径是把 Agent 放进 pipeline 内部运行,这适合批量任务(batch workloads)与评测(evaluation)。如果 Agent 本身要在生产环境持续运行,则应使用 ZenML 的 Kitaru:它把"真实发生过的运行"记录为可回放的会话,允许你在只改动一个变量(更便宜的模型、不同的提示词)的情况下对真实代码进行忠实回放,diff 两个版本后保留胜出者。回放时记录的 tool calls 直接从会话中作答,这正是其"忠实且安全"的来源。简单说:ZenML 负责 ML 流水线,Kitaru 负责 Agent。
架构示例
下面的 mermaid 图描述了本仓库 examples/deploying_agent 中的doc_analyzer流水线:CLI / curl / Web UI 请求到达 ZenML Deployment,流水线依次执行ingest_document_step→analyze_document_step→render_analysis_report_step,产物写入 Artifact Store,部署服务由 Deployer 承载:
快速开始
git clone --depth 1 https://github.com/zenml-io/zenml.git cd zenml/examples/deploying_agent uv pip install -r requirements.txt按照 examples/deploying_agent 的指引,四个步骤即可完成一个可用的 Agent 服务:
- 定义步骤:用 LLM API(OpenAI、Claude 等)构建推理步骤;
- 部署为 HTTP 服务:把 Agent 变成托管的 endpoint;
- 调用与监控:用 CLI、curl 或内置 Web UI 与 Agent 交互;
- 检查轨迹:在 ZenML 仪表盘中查看 Agent 的推理过程、工具调用与元数据。
源码级实现:doc_analyzer 流水线
流水线定义在 pipelines/doc_analyzer.py,完整展示了生产级 Agent 流水线的配置方式:
from zenml import ArtifactConfig, pipeline from zenml.config import CORSConfig, DeploymentSettings, DockerSettings docker_settings = DockerSettings( requirements="requirements.txt", environment={"OPENAI_API_KEY": "${OPENAI_API_KEY}"}, ) deployment_settings = DeploymentSettings( app_title="Document Analysis Pipeline", dashboard_files_path="ui", cors=CORSConfig(allow_origins=["*"]), ) @pipeline( settings={"docker": docker_settings, "deployment": deployment_settings}, enable_cache=False, # Disable caching for serving ) def doc_analyzer(content=None, url=None, path=None, filename=None, document_type="text"): document = ingest_document_step(content, url, path, filename, document_type) analysis = analyze_document_step(document) # OpenAI or deterministic fallback render_analysis_report_step(analysis) # HTML report for the dashboard return analysis要点解读(均有源码依据):
DockerSettings负责容器化:把requirements.txt安装进镜像,并通过${OPENAI_API_KEY}环境变量插值注入密钥,避免密钥硬编码;DeploymentSettings负责服务化:app_title设置界面标题,dashboard_files_path="ui"让ui/index.html中的 SPA 随部署自动托管,CORSConfig(allow_origins=["*"])放开跨域以便 Web 前端直连;enable_cache=False:服务化场景下关闭缓存,确保每次请求都真实执行;- 返回类型标注
ArtifactConfig(name="document_analysis", tags=["analysis", "serving"]),让分析结果以带标签的 artifact 形式落库,便于仪表盘检索。
analyze_document_step的实现位于 steps/analyze.py,其设计体现了"在线/离线双模式":
- LLM 路径:设置了
OPENAI_API_KEY时,调用 OpenAIchat.completions.create,解析结构化 JSON(摘要、关键词、情感、可读性),并记录tokens_prompt、tokens_completion、latency_ms等指标; - 回退路径:LLM 不可用时,降级为规则式分析器,用停用词过滤与词频统计提取关键词,按平均词长估算可读性分数——不依赖任何外部调用;
- 步骤内通过
try/except在两条路径间自动切换,并把analysis_method(llm或deterministic_fallback)写入metadata。
示例输出
- 自动化文档分析(见 examples/deploying_agent,下图为该示例部署后的界面截图);
- 带上下文的多轮对话机器人;
- 集成工具调用的自主工作流;
- 带检索步骤的 Agentic RAG 系统。
上图来自 deploying_agent 示例,展示了部署后的 Web 界面:支持 Direct Content(直接粘贴)、Upload File(上传文件)、URL 三种输入方式,并展示词数、处理时间、可读性分数、分析方法(Rule-Based/LLM)等指标。
相关示例
- examples/agent_outer_loop:将 ML 分类器与 Agent 结合的混合智能系统;
- examples/agentic_hitl_pipeline:在 Agent 工作流中加入动态 fan-out 与人工审批;
- examples/agent_comparison:对比不同 Agent 架构与 LLM 提供商;
- examples/agent_framework_integrations:集成 LangChain、LangGraph、LlamaIndex、CrewAI、AutoGen、Haystack 等主流 Agent 框架;
- examples/llm_finetuning:针对专项任务微调 LLM。
Path 2:构建经典机器学习流水线
使用 scikit-learn、TensorFlow、PyTorch 或其他 ML 框架构建数据处理、特征工程、训练与推理流水线。
架构示例
examples/deploying_ml_model 的客户流失预测示例将训练与推理拆成两条流水线:训练阶段generate_churn_data→train_churn_model,推理阶段predict_churn接收Customer Features(curl / SDK),产物统一写入 Artifact Store,训练与推理分别由 Orchestrator 和 Deployer 承载:
快速开始
git clone --depth 1 https://github.com/zenml-io/zenml.git cd zenml/examples/deploying_ml_model uv pip install -r requirements.txt按照 examples/deploying_ml_model 的指引操作:
- 构建流水线:数据加载 → 预处理 → 训练 → 评估;
- 部署模型:把训练好的模型作为实时 HTTP endpoint 对外服务;
- 监控性能:在仪表盘中追踪预测结果、延迟与数据漂移;
- 迭代:重训与重新部署无需改动代码——只需切换 orchestrator。
源码级实现:训练流水线
训练流水线 生成合成客户数据并训练一个随机森林分类器:
from zenml import pipeline from zenml.config import DockerSettings @pipeline( enable_cache=False, settings={"docker": DockerSettings(requirements="requirements.txt")}, ) def churn_training_pipeline( num_samples: int = 1000, test_size: float = 0.2, random_state: int = 42 ) -> Tuple[Pipeline, float]: features, target = generate_churn_data( num_samples=num_samples, random_seed=random_state ) model, accuracy = train_churn_model( features=features, target=target, test_size=test_size, random_state=random_state, ) return model, accuracy训练流水线通过Tuple[Pipeline, float]同时产出模型对象与准确率两个 artifact,其中模型会被打上production标签,供推理服务在启动时加载(详见 README 的说明)。
源码级实现:推理流水线与 Warm Container 模式
推理流水线 是 Path 2 的核心亮点,完整演示了 ZenML 的实时服务化能力:
from zenml import pipeline from zenml.config import CORSConfig, DeploymentSettings, DockerSettings from zenml.config.resource_settings import ResourceSettings @pipeline( enable_cache=False, on_init=init_model, # 部署启动时只执行一次 on_cleanup=cleanup_model, # 优雅清理 settings={ "docker": DockerSettings(requirements="requirements.txt"), "deployment": DeploymentSettings( app_title="Customer Churn Prediction Service", app_description="Real-time churn prediction with interactive web interface", app_version="1.0.0", dashboard_files_path="ui", cors=CORSConfig( allow_origins=["*"], allow_methods=["GET", "POST", "OPTIONS"], allow_headers=["*"], allow_credentials=True, ), ), "resources": ResourceSettings( memory="1GB", cpu_count=1, min_replicas=1, max_replicas=3, max_concurrency=10, ), }, ) def churn_inference_pipeline(customer_features: Dict[str, float] = {...}) -> Dict[str, Any]: return predict_churn(customer_features=customer_features)几个值得展开的配置项:
- Warm Container 模式:
on_init=init_model让模型在部署启动时一次性加载进内存,此后所有请求共享这份热模型,规避无服务器方案常见的 8~15 秒冷启动。init_model的实现见 pipelines/hooks.py,它通过Client().get_artifact_version(name_id_or_prefix="churn-model")从 Artifact Store 拉取production标签的模型版本并load()到内存,日志中还会打印加载的模型版本号;on_cleanup在服务停止时执行资源清理; ResourceSettings:控制服务的资源规格与弹性——min_replicas/max_replicas设置副本伸缩范围,max_concurrency限制并发,memory/cpu_count分配单副本资源;DeploymentSettings的CORSConfig在示例中进一步细分了allow_methods、allow_headers与allow_credentials,可精细控制跨域策略;- 推理流水线的默认入参字典完整定义了 8 个客户特征(
account_length、customer_service_calls、monthly_charges、total_charges、has_internet_service、has_phone_service、contract_length、payment_method_electronic),返回churn_probability、churn_prediction、model_version、model_status等字段。
部署命令(README 原文):
python run.py --train # 先训练,模型会被标记为 production zenml pipeline deploy pipelines.inference_pipeline.churn_inference_pipeline # 部署推理服务的 Web 界面截图如下:
上图来自 deploying_ml_model 示例:左侧填写客户特征,右侧实时展示流失概率(如 30.1%)、风险等级与预测结论。
除了自定义前端,ZenML 仪表盘还为已部署流水线内置了Playground:无需写任何代码即可在浏览器中直接向服务发送请求、查看实时预测,适合验证部署、排查问题并向团队分享可运行的示例:
上图为 deploying_ml_model 示例 中展示的 Playground 界面:左侧输入 JSON 格式的客户特征,右侧返回流失概率(如 0.018)、预测结果(0 = Will Stay)与模型状态。
你还可以通过curl直接调用/invoke接口(完整请求体见 README),并访问http://localhost:8000/docs查看自动生成的 Swagger API 文档。
示例输出
- 预测模型(回归、分类);
- 时间序列预测;
- NLP 流水线(情感分析、文本分类);
- 计算机视觉工作流;
- 模型评分与排序系统。
相关示例
- examples/e2e:包含数据验证与模型部署的端到端 ML 流水线;
- examples/e2e_nlp:领域专属的 NLP 流水线示例;
- examples/mlops_starter:带监控与治理的生产级 MLOps 配置。
Path 3:构建混合系统
在同一条流水线中结合经典 ML 模型与 AI Agent。典型用法:用分类器把请求路由到专门的 Agent,或用 Agent 增强 ML 预测结果。
架构示例
examples/agent_outer_loop 展示了"意图分类 + Agent 响应"的混合架构:客户输入进入 Agent Service,训练阶段load_data→train_classifier,服务阶段classify_intent→generate_response,全部产物与计算由统一的 Stack 承载:
快速开始
git clone --depth 1 https://github.com/zenml-io/zenml.git cd zenml/examples/agent_outer_loop uv pip install -r requirements.txt按照 examples/agent_outer_loop 的指引操作:
- 定义两个组成部分:经典 ML 分类器 step + AI Agent step;
- 串联二者:用分类器输出影响 Agent 行为(如决定调用哪个专长 Agent);
- 作为一个服务部署:整个混合系统成为单一 endpoint;
- 同时监控:在同一仪表盘中追踪 ML 指标与 Agent 轨迹。
示例输出
- 意图分类 + 专属 Agent 处理;
- 升级路径:通用 Agent → 训练分类器 → 自动路由;
- 结合多个模型与 Agent 的集成系统;
- 带验证步骤的事实核查流水线。
相关示例
- examples/agent_outer_loop:带自动意图检测的完整混合示例;
- examples/deploying_agent:从这里入手 Agent 部分;
- examples/deploying_ml_model:从这里入手 ML 部分。
三条路径通用的下一步
选定路径并跑通第一条流水线后,三条路径使用完全相同的部署模式。
远程部署:换 stack 不换代码
配置一个远程 stack(以 AWS 为例)并部署:
# 注册远程 stack(示例:AWS SageMaker + S3 + AWS Deployer) zenml stack register my-remote-stack \ --orchestrator aws-sagemaker \ --artifact-store s3-bucket \ --deployer aws # 切换 stack,你的代码一行都不用改 zenml stack set my-remote-stack以批处理模式运行:
python run.py以实时 endpoint部署:
zenml pipeline deploy pipelines.my_pipeline.my_pipeline --config deploy_config.yaml云端搭建细节参见 Deploying ZenML。Stack 切换机制更深入的解释见 Stack Components;从源码结构看,每个远程 stack 组件(如aws-sagemakerorchestrator、s3-bucketartifact store)都对应 src/zenml/integrations 下的具体集成实现,这正是"切换 stack 即切换执行环境"的底层来源。
查看仪表盘
登录并浏览流水线运行情况:
zenml login在仪表盘中你可以看到:
- Pipeline DAGs:步骤及其数据流的可视化表示;
- Artifacts:每个步骤的版本化输出(模型、报告、轨迹 traces);
- Metadata:延迟、token、指标或你自定义追踪的任何元数据;
- Timeline view:对比各步骤耗时,定位瓶颈。
核心概念回顾
无论选择哪条路径,以下五个概念贯穿始终(深入文档均已链接):
- Pipelines:编排工作流步骤并自动追踪;
- Steps:模块化、可复用的单元(数据加载、模型训练、LLM 推理等);
- Artifacts:带自动日志的版本化输出(模型、预测、轨迹、报告);
- Stacks:不修改代码即可切换执行环境(本地、远程、云端);
- Deployments:把流水线变成带内置 UI 与监控的 HTTP 服务。
至此,你应该已经理解:无论你构建的是 Agent、经典 ML 模型还是二者的混合体,ZenML 提供的都是同一条"定义步骤 → 组装流水线 → 产出 artifact → 切换 stack 部署"的主线。从 examples/quickstart 的最小骨架出发,把 examples/deploying_agent、examples/deploying_ml_model 与 examples/agent_outer_loop 三个示例跑通一遍,你就能把这条主线内化为自己的第一套可复现、可观测、可部署的 AI 流水线。
【免费下载链接】zenmlZenML 🙏: One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考