
MLOps Zoomcamp 第 4 章模型部署——Web Service、Streaming 与 Batch 三种部署方式的完整实战【免费下载链接】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 课程第 4 章Model Deployment展开完整覆盖本章的六种部署形态用 Flask Docker 打包的在线 Web Service、从 MLflow 模型仓库加载模型的 Web Service、基于 AWS Kinesis Lambda 的流式Streaming部署以及 Prefect 驱动的离线批处理Batch打分脚本。读完本文你不仅理解三种部署方式的选型逻辑还能直接复现课程配套的 Web Service 代码、MLflow 模型加载示例、Lambda 预测函数 与 批量打分流程并掌握其 Docker 镜像构建、模型仓库寻址、Kinesis 事件解码与 Prefect 定时部署等关键细节。4.1 模型部署的三种方式先建立选型视角本章的开篇04-deployment/README.md 的 4.1 节指出生产环境中的模型部署方式可以从在线Online/离线Offline两个维度来区分落到工程实现上主要是三种部署方式形态典型技术栈适用场景课程对应小节Web Service在线、同步请求-响应Flask Docker低延迟的实时打分4.2 / 4.3Streaming在线、事件驱动AWS Kinesis Lambda事件流持续到达、按条处理4.4Batch离线、批量打分脚本 Prefect 调度定时对历史/整月数据打分4.5 / 4.6三种方式都服务于同一个课程案例纽约出租车行程时长预测ride duration prediction训练产物是DictVectorizer 回归模型的 scikit-learn pipelineRandom Forest 或线性回归。理解了输入特征PULocationID、DOLocationID、trip_distance和输出duration单位分钟就可以把三种部署方式当作同一个 predict 函数在不同调用形态下的封装来看待这正是本章全部代码的公共骨架。4.2 Web Service用 Flask 和 Docker 部署模型工程结构四步走web-service 模块的 README 把该节的工作拆成四步用 Pipenv 创建虚拟环境Pipfile声明依赖编写预测脚本predict.py把脚本包进一个 Flask App把 App 打包成 Docker 镜像。预测服务predict.py 逐段解析04-deployment/web-service/predict.py 是本章最核心的服务文件其结构可以拆成四部分import pickle from flask import Flask, request, jsonify with open(lin_reg.bin, rb) as f_in: (dv, model) pickle.load(f_in)服务启动时一次性从lin_reg.bin反序列化出DictVectorizerdv和模型model并常驻内存——这是 Web Service 的低延迟来源避免每次请求都重新加载模型。def prepare_features(ride): features {} features[PU_DO] %s_%s % (ride[PULocationID], ride[DOLocationID]) features[trip_distance] ride[trip_distance] return features def predict(features): X dv.transform(features) preds model.predict(X) return float(preds[0])prepare_features把原始请求里的上下车位置编码成类别特征PU_DO形如43_151再拼上数值特征trip_distance与训练时prepare_dictionaries的口径严格一致。predict用dv.transform做向量化后交给模型打分并取float(preds[0])保证返回 JSON 可序列化的标量。app Flask(duration-prediction) app.route(/predict, methods[POST]) def predict_endpoint(): ride request.get_json() features prepare_features(ride) pred predict(features) result {duration: pred} return jsonify(result) if __name__ __main__: app.run(debugTrue, host0.0.0.0, port9696)只暴露一个POST /predict端点if __name__ __main__分支仅用于本地开发调试生产路径走 Docker 里的 gunicorn见下。请求体由 test.py 给出可直接用来验证服务import requests ride { PULocationID: 10, DOLocationID: 50, trip_distance: 40 } url http://localhost:9696/predict response requests.post(url, jsonride) print(response.json())依赖声明与 Docker 打包Pipfile 中运行时依赖为scikit-learn1.0.2、flask、gunicornrequests只在[dev-packages]中用于本地测试python_version 3.9。Dockerfile 展示了先装依赖再拷代码的分层写法能最大化利用构建缓存FROM python:3.9.7-slim RUN pip install -U pip RUN pip install pipenv WORKDIR /app COPY [ Pipfile, Pipfile.lock, ./ ] RUN pipenv install --system --deploy COPY [ predict.py, lin_reg.bin, ./ ] EXPOSE 9696 ENTRYPOINT [ gunicorn, --bind0.0.0.0:9696, predict:app ]几个要点pipenv install --system --deploy--system直接装入镜像 Python 环境--deploy保证Pipfile.lock与Pipfile一致否则构建失败——这是把锁文件作为可复现构建门禁的用法lin_reg.bin随代码一起COPY进镜像模型与代码同生命周期入口命令不是 Flask 自带的app.run而是gunicorn --bind0.0.0.0:9696 predict:app用 WSGI 服务器承载生产流量EXPOSE 9696声明服务端口。构建与运行命令来自模块 README可直接复制执行docker build -t ride-duration-prediction-service:v1 .docker run -it --rm -p 9696:9696 ride-duration-prediction-service:v1构建后执行 test.py即可获得形如{duration: ...}的响应。4.3 Web Service 进阶从 MLflow 模型仓库加载模型4.2 节的模型是烧进镜像的 pickle 文件模型升级就要重新构建镜像。web-service-mlflow 模块演示了更工程化的方式镜像不携带模型文件运行时从 MLflow 模型仓库Tracking Server / S3按 run_id 拉取模型。训练并注册把模型放进 pipelinerandom-forest.ipynb 展示了训练侧做法读取 parquet 数据、计算duration上下车时间差换算成分钟并过滤 1~60 分钟样本、把PULocationID/DOLocationID拼接为PU_DO然后在一个 MLflow run 里完成参数、指标、模型的记录with mlflow.start_run(): params dict(max_depth20, n_estimators100, min_samples_leaf10, random_state0) mlflow.log_params(params) pipeline make_pipeline( DictVectorizer(), RandomForestRegressor(**params, n_jobs-1) ) pipeline.fit(dict_train, y_train) y_pred pipeline.predict(dict_val) rmse mean_squared_error(y_pred, y_val, squaredFalse) mlflow.log_metric(rmse, rmse) mlflow.sklearn.log_model(pipeline, artifact_pathmodel)关键点mlflow.sklearn.log_model(pipeline, ...)把整个 pipeline向量化器 模型作为pyfunc 模型记录这样下游只需mlflow.pyfunc.load_model就能拿到端到端可预测对象不需要分别管理 DictVectorizer 和 RandomForestRegressor。启动带 S3 后端的 MLflow Server模块 README 给出了 tracking server 的启动方式mlflow server \ --backend-store-urisqlite:///mlflow.db \ --default-artifact-roots3://mlflow-models-alexey/--backend-store-uri指向元数据库这里是本地 SQLite--default-artifact-root把模型产物默认存到 S3 桶从而把元数据与大文件解耦。不依赖 Tracking Server 直接加载模型load_model.ipynb 演示了两种寻址方式# 方式一走 tracking server 的 runs:/ URI TRACKING_URL http://127.0.0.1:5000 mlflow.set_tracking_uri(TRACKING_URL) run_id 6dd459b11b4e48dc862f4e1019d166f6 logged_model fruns:/{run_id}/model loaded_model mlflow.pyfunc.load_model(logged_model) # 方式二直接走 S3 artifact 路径无需 tracking server loaded_model mlflow.pyfunc.load_model( s3://mlflow-models-alexey/1/6dd459b11b4e48dc862f4e1019d166f6/artifacts/model/)随后用与 Web Service 相同口径的特征直接预测request { lpep_pickup_datetime: 2021-01-01 00:15:56, PULocationID: 43, DOLocationID: 151, passenger_count: 1.0, trip_distance: 1.01 } features {} features[PU_DO] %s_%s % (request[PULocationID], request[DOLocationID]) features[trip_distance] request[trip_distance] loaded_model.predict(features)这里能看出一个设计细节MLflow 的 pyfunc 模型接口接受与训练时相同形态的 dict 输入PU_DOtrip_distance把特征工程封装在 pipeline 内部让训练口径和部署口径天然一致。此外README 还给出了 CLI 下载 artifact 的完整命令可用于把模型物化到本地再打进镜像export MLFLOW_TRACKING_URIhttp://127.0.0.1:5000 export MODEL_RUN_ID6dd459b11b4e48dc862f4e1019d166f6 mlflow artifacts download \ --run-id ${MODEL_RUN_ID} \ --artifact-path model \ --dst-path .服务端的两种部署形态web-service-mlflow/predict.py 与 4.2 版的差异集中在两处import os import mlflow RUN_ID os.getenv(RUN_ID) logged_model fs3://mlflow-models-alexey/1/{RUN_ID}/artifacts/model # logged_model fruns:/{RUN_ID}/model model mlflow.pyfunc.load_model(logged_model)模型地址由环境变量RUN_ID参数化走 S3 直连默认或走 tracking serverruns:/形式源码中留有注释切换换模型只需换环境变量不需要重构建镜像响应体里额外返回model_version即RUN_ID便于下游审计这条预测来自哪个模型版本result { duration: pred, model_version: RUN_ID }另外从 random-forest.ipynb 可以看到若只需向量化器而非完整模型也可以用MlflowClient.download_artifacts(run_idRUN_ID, pathdict_vectorizer.bin)单独下载 artifact 后用pickle.load加载——这是按 artifact 粒度取物的补充手段。4.4 Streaming用 Kinesis 和 Lambda 部署模型可选但推荐4.4 节的场景是事件流行程事件持续写入 Kinesis 流Lambda 函数被流触发逐条解码、打分并把预测结果再写回一条 Kinesis 流。模块 README 说明因涉及 AWS 付费服务本视频是可选跟练但课程第 6 章Best Practices的基础正是这条流水线所以仍建议观看。Lambda 函数lambda_function.py 的完整逻辑04-deployment/streaming/lambda_function.py 的核心流程import os import json import base64 import boto3 import mlflow kinesis_client boto3.client(kinesis) PREDICTIONS_STREAM_NAME os.getenv(PREDICTIONS_STREAM_NAME, ride_predictions) RUN_ID os.getenv(RUN_ID) logged_model fs3://mlflow-models-alexey/1/{RUN_ID}/artifacts/model model mlflow.pyfunc.load_model(logged_model) TEST_RUN os.getenv(TEST_RUN, False) True与 4.3 相同的模型寻址方式模块级加载一次模型靠RUN_ID环境变量指定版本结果输出到PREDICTIONS_STREAM_NAME默认ride_predictions指定的输出流。TEST_RUN环境变量用于本地测试时跳过真正写流。def lambda_handler(event, context): predictions_events [] for record in event[Records]: encoded_data record[kinesis][data] decoded_data base64.b64decode(encoded_data).decode(utf-8) ride_event json.loads(decoded_data) ride ride_event[ride] ride_id ride_event[ride_id] features prepare_features(ride) prediction predict(features) prediction_event { model: ride_duration_prediction_model, version: 123, prediction: { ride_duration: prediction, ride_id: ride_id } } if not TEST_RUN: kinesis_client.put_record( StreamNamePREDICTIONS_STREAM_NAME, Datajson.dumps(prediction_event), PartitionKeystr(ride_id) ) predictions_events.append(prediction_event) return {predictions: predictions_events}几个值得注意的实现细节Kinesis 记录数据是 base64 编码的必须base64.b64decode(encoded_data).decode(utf-8)后再json.loads——这是所有 Kinesis 消费者都要处理的一步一次触发事件可能包含多条Records函数对每条独立打分PartitionKeystr(ride_id)让同一行程的结果落在同一分区保证顺序性输出事件结构固定为{model, version, prediction: {ride_duration, ride_id}}是下游消费方的契约。输入数据格式与发送记录事件 JSON 的结构见 streaming/README.md 与 test.py{ ride: { PULocationID: 130, DOLocationID: 205, trip_distance: 3.66 }, ride_id: 123 }用 AWS CLI 向流ride_events发送一条记录KINESIS_STREAM_INPUTride_events aws kinesis put-record \ --stream-name ${KINESIS_STREAM_INPUT} \ --partition-key 1 \ --data { ride: { PULocationID: 130, DOLocationID: 205, trip_distance: 3.66 }, ride_id: 156 }本地测试与 Lambda 容器镜像streaming/test.py 把一份真实的 Kinesis 事件含 base64data字段、eventSourceARN等元信息直接喂给lambda_handler在不上云的情况下验证整条解码 → 打分 → 组事件的链路import lambda_function result lambda_function.lambda_handler(event, None) print(result)streaming/Dockerfile 用于docker-lambda本地执行 Lambda 镜像FROM public.ecr.aws/lambda/python:3.9 RUN pip install -U pip RUN pip install pipenv COPY [ Pipfile, Pipfile.lock, ./ ] RUN pipenv install --system --deploy COPY [ lambda_function.py, ./ ] CMD [ lambda_function.lambda_handler ]注意入口是CMD [lambda_function.lambda_handler]——Lambda 容器镜像要求CMD指向模块.处理函数与 Web Service 的 gunicorn 入口形成对照同一份加载模型 predict逻辑换一层事件驱动外壳即可在流式场景复用。4.5 Batch把训练 Notebook 改造成打分脚本Batch 部署的目标是定时对整月历史数据打分并把预测结果落盘归档。batch 模块 README 概括了三个改造步骤把训练 notebook 变成应用模型的 notebook再把 notebook 变成脚本最后清理并参数化。打分脚本score.py04-deployment/batch/score.py 是参数化的 Prefect 2.0 flow关键函数逐一说明generate_uuids(n)给每行生成ride_id保证预测结果可回溯、可去重read_dataframe(filename)读取 S3 上的月度 parquet如s3://nyc-tlc/trip data/green_tripdata_2021-03.parquet同样计算duration并过滤 1~60 分钟样本prepare_dictionaries(df)与训练侧一致的PU_DOtrip_distance口径load_model(run_id)从s3://mlflow-models-alexey/1/{run_id}/artifacts/model用mlflow.pyfunc.load_model加载模型——与 4.3/4.4 完全一致的寻址方式save_results(...)输出 parquet 含ride_id、上下车时间与位置、actual_duration、predicted_duration、diff以及model_version即 run_id形成版本可追溯的预测台账df_result[actual_duration] df[duration] df_result[predicted_duration] y_pred df_result[diff] df_result[actual_duration] - df_result[predicted_duration] df_result[model_version] run_id df_result.to_parquet(output_file, indexFalse)路径参数化在get_paths中完成输入固定为运行日前一个月的 tripdata输出按taxi_type/year/month/run_id分层写入s3://nyc-duration-prediction-alexey/def get_paths(run_date, taxi_type, run_id): prev_month run_date - relativedelta(months1) input_file fs3://nyc-tlc/trip data/{taxi_type}_tripdata_{year:04d}-{month:02d}.parquet output_file (fs3://nyc-duration-prediction-alexey/ ftaxi_type{taxi_type}/year{year:04d}/month{month:02d}/{run_id}.parquet) return input_file, output_fileCLI 入口run()从sys.argv读取四个参数并触发 flowpython score.py green 2021 3 e1efc53e9bd149078b0c12aeaa6365df即taxi_typegreen、year2021、month3打分 2021 年 2 月数据、run_id指定模型版本。定时部署与历史回填score_deploy.py 把该 flow 注册为 Prefect Deployment用 Cron 表达式实现每月 2 日凌晨 3 点跑上月数据deployment Deployment.build_from_flow( flowride_duration_prediction, nameride_duration_prediction, parameters{ taxi_type: green, run_id: e1efc53e9bd149078b0c12aeaa6365df, }, scheduleCronSchedule(cron0 3 2 * *), work_queue_nameml, ) deployment.apply()score_backfill.py 则演示历史回填用relativedelta(months1)从 2021-03 循环到 2022-04逐月调用同一个 flow。这两段代码共同说明了 Batch 部署的完整生命周期脚本化 → 定时化 → 可回填。依赖方面batch/Pipfile 声明了scikit-learn1.0.2、prefect2.0b6、mlflow、pandas、boto3、pyarrow、s3fs其中s3fs/pyarrow支撑pd.read_parquet直接读写 S3。4.6 Batch 打分接入 Mage无视频思路复用课程原文对 4.6 节只给出三句话你已经会了Connect to MLflowCreate a transformation blockGet the model from the registry, apply it。对照 4.5 节的代码可以精确翻译这三步mlflow.set_tracking_uri/mlflow.pyfunc.load_model对应 1、3以及把prepare_dictionaries model.predict包成一个转换块对应 2。也就是说无论调度器是 Prefect 还是 MageBatch 打分的核心代码都是score.py里apply_modeltask 的同一套逻辑。作业与延伸阅读本章作业见 cohorts/2025/04-deployment/homework.md要求基于给定数据完成部署实战04-deployment/README.md的 Notes 区还汇总了多位学员的部署笔记含 GCP 路线可作横向参考。小结一张表收束本章小节部署方式模型加载方式入口形态关键文件4.2Web Servicepickle 文件随镜像 COPYgunicorn Flask/predictweb-service/Dockerfile4.3Web Service MLflow运行时按RUN_ID从 S3/tracking server 拉取同上响应含model_versionweb-service-mlflow/predict.py4.4Streaming同 4.3Lambda handlerKinesis base64 解码streaming/lambda_function.py4.5Batch同 4.3Prefect flow Cron Deployment 回填脚本batch/score.py从源码结构看三种部署形态共享同一个契约输入是{PU_DO, trip_distance}形态的特征 dict模型统一由 MLflow pyfunc 接口加载且以run_id作为版本标识输出统一标注model_version。抓住这个契约就抓住了 MLOps Zoomcamp 第 4 章的部署方法论也为第 5 章监控与第 6 章基于 Kinesis/Lambda 流水线的最佳实践打下了直接的基础。【免费下载链接】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),仅供参考