ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

MLOps Zoomcamp 2023 Orchestration 篇:用 Prefect 2.0 构建可观测的 ML 训练流水线

MLOps Zoomcamp 2023 Orchestration 篇:用 Prefect 2.0 构建可观测的 ML 训练流水线 MLOps Zoomcamp 2023 Orchestration 篇用 Prefect 2.0 构建可观测的 ML 训练流水线【免费下载链接】mlops-zoomcampFree MLOps course from DataTalks.Club. Register here to get notified about the next cohort项目地址: https://gitcode.com/GitHub_Trending/ml/mlops-zoomcamp本文基于 MLOps Zoomcamp 2023 队列第 3 章Orchestration and ML Pipelines的配套文档与代码目录展开讲解如何用 Prefect 2.0 把一个普通的 Python 模型训练脚本改造成可编排、可观测、可部署的机器学习工作流。读完后你将掌握flow/task的核心用法、任务重试机制、将流水线部署为 Prefect Deployment 的操作流程以及通过prefect-aws的 S3 数据块与 Markdown 工件实现“云端取数 实验报告回传”的完整实战方案。课程模块结构与学习路线该章节的入口文档是 cohorts/2023/03-orchestration/prefect/README.md其章节规划同步记录在 cohorts/2023/03-orchestration/prefect/meta.json 中。模块编号为 3标题为Orchestration and ML Pipelines共 7 个单元单元主题说明3.1Introduction to Workflow Orchestration工作流编排概念入门3.2Introduction to Prefect认识 Prefectflow 与 task3.3Prefect Workflow把训练流水线改造成 Prefect 工作流3.4Deploying Your Workflow将工作流部署为 Deployment3.5Working with Deployments部署的管理、S3 数据块与工件artifacts3.6Prefect Cloud可选使用托管版 Prefect Cloud3.7Homework本章作业仓库中每个视频单元对应一个同名代码目录3.2/3.6/代码按单元递进从最小的“猫狗冷知识”演示到纽约出租车行程时长预测的完整 XGBoost 训练流水线再到 S3 取数与报告工件。这种“先玩具级示例、后真实 ML 流水线”的组织方式使得每一小节都能独立运行、独立验证。环境准备与快速搭建Quick Setup原文档的 Quick setup 部分给出了完整的最小可用环境这里完整继承并结合 requirements.txt 的锁定版本做补充。1. 安装依赖在 Python 3.10.12或相近版本的 conda 环境中执行pip install -r requirements.txt其中与本章直接相关的锁定版本为prefect2.10.8 prefect-aws0.3.1 mlflow2.3.1 xgboost1.7.5 pandas2.0.1 scikit_learn1.2.2 fastparquet2023.4.0这里prefect2.10.8明确了本章全部代码面向Prefect 2.x而非 1.x 或 3.xflow/task装饰器、prefect deploy部署模型、S3Bucket数据块block等 API 都是 2.x 时代的形态。prefect-aws则提供 3.5 节使用的AwsCredentials与S3Bucketblock。2. 本地启动 Prefect 服务器另开一个终端并激活同一 conda 环境然后prefect server start该命令在本机启动 Prefect API UI 服务flow 的运行记录、部署列表都会上报到这里。本地自托管是最常见的学习与内网生产方式无需注册任何账号。3. 可选改用 Prefect Cloud不想自己维护 server 时可以使用托管版 Prefect Cloud原文档标注为 optional在 3.6 单元专门讲解。在终端中认证prefect cloud login认证后可以通过 Prefect 的profile机制在自托管服务器与 Cloud 之间切换即同一份代码既能对本地prefect server start上报也能对 Cloud 上报无需改动流水线本身。3.2 认识 Prefect从 flow 与 task 讲起理解 Prefect 的最小载体是两个装饰器flow把普通函数升级为一条工作流task把其中可独立观测、可重试的单元升级为任务。仓库在 3.2/cat_dog_facts.py 中给出了一个完整的两级嵌套示例import httpx from prefect import flow flow def fetch_cat_fact(): A flow that gets a cat fact return httpx.get(https://catfact.ninja/fact?max_length140).json()[fact] flow def fetch_dog_fact(): A flow that gets a dog fact return httpx.get( https://dogapi.dog/api/v2/facts, headers{accept: application/json}, ).json()[data][0][attributes][body] flow(log_printsTrue) def animal_facts(): cat_fact fetch_cat_fact() dog_fact fetch_dog_fact() print(f: {cat_fact} \n: {dog_fact}) if __name__ __main__: animal_facts()从源码结构看这里有一个值得注意的细节animal_facts是一个 flow而它调用的fetch_cat_fact、fetch_dog_fact同样是 flow 而不是 task——即 flow 可以嵌套调用 flow形成子工作流subflow。log_printsTrue参数则让任务体内部的print输出自动汇入 Prefect 日志系统而不是散落在标准输出里。第二个文件 3.2/cat_facts.py 演示了task与重试机制import httpx from prefect import flow, task task(retries4, retry_delay_seconds0.1, log_printsTrue) def fetch_cat_fact(): cat_fact httpx.get(https://f3-vyx5c2hfpq-ue.a.run.app/) #An endpoint that is designed to fail sporadically if cat_fact.status_code 400: raise Exception() print(cat_fact.text) flow def fetch(): fetch_cat_fact() if __name__ __main__: fetch()该任务的三个参数含义清晰retries4失败后最多再重试 4 次retry_delay_seconds0.1每次重试前等待 0.1 秒被调用的接口是课程中“被设计成间歇性失败”的端点因此该示例在 UI 中会稳定呈现“失败 → 自动重试 → 最终成功”的任务运行轨迹非常适合直观观察 Prefect 对瞬时故障的容错能力。对 ML 场景而言重试参数通常加在读数据这类依赖外部存储S3、数据库的环节而本章 3.3 起的训练流水线正是这样做的见下文read_data。3.3 把训练流水线改造成 Prefect Workflow本单元的代码目录 3.3/ 中有一对“前后对照”文件是理解 Prefect 改造成本的最佳素材orchestrate_pre_prefect.py纯过程式脚本无任何 Prefect 依赖orchestrate.py同一逻辑加上flow/task之后的 Prefect 工作流版本。两者的函数体几乎逐行相同差异只在于装饰器与导入。对照两版源码可以准确回答“接入 Prefect 到底改了多少代码”业务逻辑零改动仅增加 3 个装饰器、2 个导入、以及个别 task 上的重试参数。改造后的流水线结构Prefect 版的main_flow由三个 task 串联构成“Load → Transform → Train”三段式from prefect import flow, task task(retries3, retry_delay_seconds2) def read_data(filename: str) - pd.DataFrame: Read data into DataFrame df pd.read_parquet(filename) # ... 时间列转 datetime、计算 duration 并过滤 [1, 60] 分钟 ... task def add_features(df_train, df_val): Add features to the model # 构造 PU_DO 交叉特征DictVectorizer 向量化返回 X/y 与 dv task(log_printsTrue) def train_best_model(X_train, X_val, y_train, y_val, dv) - None: train a model with best hyperparams and write everything out # mlflow 记录参数/指标/工件训练 XGBoost 回归模型 flow def main_flow( train_path: str ./data/green_tripdata_2021-01.parquet, val_path: str ./data/green_tripdata_2021-02.parquet, ) - None: The main training pipeline # MLflow settings mlflow.set_tracking_uri(sqlite:///mlflow.db) mlflow.set_experiment(nyc-taxi-experiment) # Load df_train read_data(train_path) df_val read_data(val_path) # Transform X_train, X_val, y_train, y_val, dv add_features(df_train, df_val) # Train train_best_model(X_train, X_val, y_train, y_val, dv)以上摘录自 3.3/orchestrate.py业务细节与 pre-Prefect 版本完全一致数据任务read_data读取纽约出租车 Green 行程的 parquet 文件默认train_path为 2021-01、val_path为 2021-02把上车/下车时间转成 datetime计算以分钟计的duration并只保留 160 分钟的行程PULocationID/DOLocationID统一转为字符串。它被标记为task(retries3, retry_delay_seconds2)——读数据是典型的外部依赖环节失败如文件下载不完整后间隔 2 秒重试 3 次。特征任务add_features构造上车-下车位置组合特征PU_DO用DictVectorizer将分类特征加trip_distance向量化为稀疏矩阵同时返回向量化器dv本身以便后续持久化保证训练/推理预处理一致。训练任务train_best_model在mlflow.start_run()上下文中用固定的best_paramslearning_rate≈0.0959、max_depth30、objectivereg:linear、seed42等训练 XGBoost 回归模型num_boost_round100且early_stopping_rounds20计算验证集 RMSE 并写入 MLflow将DictVectorizerpickle 到models/preprocessor.b并以mlflow.log_artifact记录为preprocessor工件最后用mlflow.xgboost.log_model注册模型。MLflow 追踪地址为本地 SQLitesqlite:///mlflow.db实验名为nyc-taxi-experiment。值得注意的是任务间的数据传递方式add_features的返回签名显式标注了csr_matrix、np.ndarray与DictVectorizer类型3.3/orchestrate.py。从源码结构看这种“显式返回类型 显式入参类型”的写法让 Prefect 在任务间传递结果时类型清晰也为 UI 中查看每个 task 的输入输出提供了可读性。3.4 Deploying Your Workflow从本地运行到部署3.4 单元的 3.4/orchestrate.py 与 3.3 版本保持同一份流水线代码区别在于运行方式不再以python orchestrate.py直接执行而是把main_flow**部署deploy**为 Prefect 的 Deployment 对象。结合仓库中 Prefect 锁定版本2.10.8该单元的部署流程为# 1. 在流水线文件所在目录从 flow 函数生成并推送部署 prefect deploy orchestrate.py执行后Prefect CLI 会交互式询问部署名称、工作区work pool等并自动注册一个带默认参数的 Deployment——这些默认参数正是main_flow函数签名中的train_path、val_path缺省值。部署完成后即可# 2. 查看已有部署 prefect deployment ls # 3. 通过 API 触发一次运行可覆盖 flow 入参 prefect deployment run deployment-name --param train_path./data/green_tripdata_2021-01.parquet部署模型带来的核心变化是flow 的运行由“人工敲 python 脚本”变为“由 Prefect Server 管理的一次次 Run”每次运行的任务状态、日志、重试记录都在 UI 中可观测、可回放。这也是 3.4 与 3.3 使用相同代码却能区分成两节的原因——编排价值的兑现点在部署而不是装饰器本身。3.5 Working with DeploymentsS3 数据块与实验报告工件3.5 单元3.5/解决两个实际问题训练数据不在本地而存在 S3 上、训练结果需要一份可视化报告回到 Prefect UI。目录中包含三个文件create_s3_bucket_block.py一次性创建并保存 AWS 凭据与 S3 桶两个 Prefectblockorchestrate.py与 3.4 相同的本地版流水线保留作对照orchestrate_s3.py真正使用 S3 取数并产出 Markdown 工件的增强版。第一步保存 S3 数据块create_s3_bucket_block.py完整代码来自 3.5/create_s3_bucket_block.pyfrom time import sleep from prefect_aws import S3Bucket, AwsCredentials def create_aws_creds_block(): my_aws_creds_obj AwsCredentials( aws_access_key_id123abc, aws_secret_access_keyabc123 ) my_aws_creds_obj.save(namemy-aws-creds, overwriteTrue) def create_s3_bucket_block(): aws_creds AwsCredentials.load(my-aws-creds) my_s3_bucket_obj S3Bucket( bucket_namemy-first-bucket-abc, credentialsaws_creds ) my_s3_bucket_obj.save(names3-bucket-example, overwriteTrue) if __name__ __main__: create_aws_creds_block() sleep(5) create_s3_bucket_block()要点block 的save(name..., overwriteTrue)会持久化到 Prefect Server之后任何环境都可用AwsCredentials.load(my-aws-creds)、S3Bucket.load(...)按名取回——凭据与桶配置从环境变量/代码中解耦交给编排平台统一管理。文件中的 key 是占位符实际使用时需替换为真实的 AWS 凭证与桶名sleep(5)是避免两个 block 在 UI 中同时写入造成前端状态不同步的小技巧。第二步流水线中加载数据块并下载数据增强版 orchestrate_s3.py 相对 3.4 版本只有两处增量from prefect_aws import S3Bucket from prefect.artifacts import create_markdown_artifact from datetime import date flow def main_flow_s3( train_path: str ./data/green_tripdata_2021-01.parquet, val_path: str ./data/green_tripdata_2021-02.parquet, ) - None: # MLflow settings mlflow.set_tracking_uri(sqlite:///mlflow.db) mlflow.set_experiment(nyc-taxi-experiment) # Load s3_bucket_block S3Bucket.load(s3-bucket-block) s3_bucket_block.download_folder_to_path(from_folderdata, to_folderdata) df_train read_data(train_path) df_val read_data(val_path) # ... 后续 Transform / Train 与本地版完全一致 ...S3Bucket.load(s3-bucket-block)按名加载块download_folder_to_path(from_folderdata, to_folderdata)把 S3 上data/前缀下的 parquet 文件整批拉到本地data/目录使后续read_data的本地路径逻辑完全不用改动。这里“下载整文件夹后走本地读数据”的写法也解释了为什么read_data依然保留retries3, retry_delay_seconds2网络取数正是最需要自动重试的一环。第三步把 RMSE 报告作为 Markdown 工件回传train_best_model任务在记录 MLflow 指标之后额外生成一份 Markdown 报告并通过create_markdown_artifact发布到 Prefect UI见 orchestrate_s3.pymarkdown__rmse_report f# RMSE Report ## Summary Duration Prediction ## RMSE XGBoost Model | Region | RMSE | |:----------|-------:| | {date.today()} | {rmse:.2f} | create_markdown_artifact( keyduration-model-report, markdownmarkdown__rmse_report )工件artifact是 Prefect 2.x 将“运行产出”与“运行日志”分离的机制报告以duration-model-report为 key 挂在对应 flow run 的 Artifacts 页签上打开运行详情即可看到渲染后的表格无需再去翻 MLflow。3.6 单元切到 Prefect Cloud 的增量变化3.6/orchestrate_s3.py 与 3.6/create_s3_bucket_block.py 的流水线代码与 3.5 基本一致差异仅在于报告 Markdown 的缩进排版。从源码结构看这体现了本章的设计意图切换到 Prefect Cloud 时流水线代码几乎不需要动主要操作集中在前面 Quick setup 中提到的prefect cloud login认证以及用 profile 把客户端指向 Cloud。换言之block、deployment、artifact 的 API 在自托管与托管两种后端下是同构的这正是 Prefect 2.x “代码不变、后端可换”的部署灵活性。作业与延伸阅读本章作业Homework单元 3.7见 cohorts/2023/03-orchestration/homework.md。上一届2022对应章节使用了 Prefect 1.x 时代的写法含prefect deploy的前身流程可作为版本演进的对照阅读cohorts/2022/03-orchestration/README.md。版本与适用前提说明结合 requirements.txt 与代码实际调用的 API使用本章材料时需要注意以下前提Prefect 版本代码基于prefect2.10.8flow/task、prefect deploy、S3Bucketblock、create_markdown_artifact均为 2.x API在 Prefect 3.x 中部分入口如部分 CLI 与默认工作池行为已调整直接照搬本目录代码到 3.x 前应先核对官方迁移说明。Python 环境文档推荐 conda Python 3.10.12 或相近版本依赖pandas 2.0.1、xgboost 1.7.5、mlflow 2.3.1 等按锁定版本安装以避免 API 漂移。两种后端二选一本地prefect server start或prefect cloud login使用 Cloud通过 profile 切换block 与 deployment 均存储在所选后端中切换前需确认 block 已存在于对应后端。S3 相关步骤create_s3_bucket_block.py中的密钥与桶名为占位符实际运行需替换为有效 AWS 凭据并保证该凭据对目标桶具备读写权限orchestrate_s3.py 依赖名为s3-bucket-block的块已按文档流程保存。数据前提main_flow默认读取./data/green_tripdata_2021-01.parquet与./data/green_tripdata_2021-02.parquet或从 S3 下载到本地data/目录运行前需确保这两个 parquet 文件存在配套的数据探索 notebook如 3.3/duration_prediction_explore.ipynb可用于理解数据字段与特征构造动机。小结本章以纽约出租车行程时长预测为实战载体展示了 Prefect 2.0 接入既有 ML 脚本的完整路径——task化拆分带来任务级观测与重试prefect deploy把脚本变成平台管理的 Deploymentprefect-aws的 block 让云端数据源即插即用Markdown artifact 则把实验报告直接带回运行详情页。全程业务代码近乎零改动这是 Prefect 类编排框架在 MLOps 场景中的典型落地方式。【免费下载链接】mlops-zoomcampFree MLOps course from DataTalks.Club. Register here to get notified about the next cohort项目地址: https://gitcode.com/GitHub_Trending/ml/mlops-zoomcamp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表