
MLOps Zoomcamp用 Amazon 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本文基于 MLOps Zoomcamp 第 4 章Model Deployment中“Streaming”小节的配套文档 streaming/README.md 展开系统讲解如何用 Amazon Kinesis 流 AWS Lambda MLflow 模型构建一条流式Streaming推理管道从创建 IAM 角色、编写 Lambda 函数、本地测试到创建输入/输出 Kinesis 流、发送记录、Docker 打包、ECR 镜像发布并附源码级实现分析与后续第 6 章的生产级 Terraform 演进路径。读完本文你可以完整复刻课程视频 4.4 中的实验并理解事件结构、base64 编解码与测试开关等关键机制。1. 背景三种部署模式中的流式推理在 04-deployment/README.md 中课程将模型部署分为三种方式Web-service在线服务Flask Docker 部署并进一步用 MLflow 模型仓库加载模型Streaming流式推理Kinesis Lambda本文的主题Batch批量打分定时对历史数据批量打分。其中 Streaming 小节4.4在课程 README 中被标注为可选项原因是它依赖会产生费用的 AWS 服务但文档同时指出第 6 章Best Practices的内容正是基于本节内容构建的因此强烈建议学习。这一演进关系在本仓库中可以得到印证06-best-practices/code 目录中包含了完整的 Terraform 基础设施、Docker 镜像打包与 Kinesis 集成测试本质上是把第 4 章的手工流式管道工程化、可测试化。1.1 整体流程streaming/README.md 开头列出的 6 步即本节的实施路线Scenario确定场景——用 Lambda 消费 Kinesis 上的打车行程事件输出预测结果到另一条流Creating the role创建 IAM 角色授予 Lambda 读取 Kinesis、写入输出流的权限Create a Lambda function, test it创建 Lambda 函数并先行测试Create a Kinesis stream创建输入流ride_events以及输出流ride_predictionsConnect the function to the stream建立事件源映射让 Lambda 订阅输入流Send the records向流中写入测试记录并验证预测输出。README 还引用了 AWS 官方的 “Using Amazon Lambda with Amazon Kinesis” 教程作为外部参考见 README 中的 Links 部分。2. 事件数据结构Kinesis 记录与 base64 编码理解整条管道的前提是理解 Lambda 收到的事件长什么样。2.1 输入记录ride eventREADME 中给出的业务记录示例是一个 JSON 对象包含ride行程特征和ride_id两个字段{ ride: { PULocationID: 130, DOLocationID: 205, trip_distance: 3.66 }, ride_id: 123 }发送该记录到 Kinesis 的命令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 }最基础的冒烟测试也可以只发一段纯文本KINESIS_STREAM_INPUTride_events aws kinesis put-record \ --stream-name ${KINESIS_STREAM_INPUT} \ --partition-key 1 \ --data Hello, this is a test.2.2 关键点data 字段是 base64 编码Kinesis 的put-record会把--data内容 base64 编码后存储因此 Lambda 侧拿到的record[kinesis][data]是编码字符串需要解码还原。README 给出了标准解码方式base64.b64decode(data_encoded).decode(utf-8)2.3 完整的测试事件Test event在 Lambda 控制台测试时使用的 event 结构完整继承自真实触发场景。README 给出的测试事件如下其中账号 ID 与区域为课程作者的示例值{ Records: [ { kinesis: { kinesisSchemaVersion: 1.0, partitionKey: 1, sequenceNumber: 49630081666084879290581185630324770398608704880802529282, data: ewogICAgICAgICJyaWRlIjogewogICAgICAgICAgICAiUFVMb2NhdGlvbklEIjogMTMwLAogICAgICAgICAgICAiRE9Mb2NhdGlvbklEIjogMjA1LAogICAgICAgICAgICAidHJpcF9kaXN0YW5jZSI6IDMuNjYKICAgICAgICB9LCAKICAgICAgICAicmlkZV9pZCI6IDI1NgogICAgfQ, approximateArrivalTimestamp: 1654161514.132 }, eventSource: aws:kinesis, eventVersion: 1.0, eventID: shardId-000000000000:49630081666084879290581185630324770398608704880802529282, eventName: aws:kinesis:record, invokeIdentityArn: arn:aws:iam::XXXXXXXXX:role/lambda-kinesis-role, awsRegion: eu-west-1, eventSourceARN: arn:aws:kinesis:eu-west-1:XXXXXXXXX:stream/ride_events } ] }结构上有三层顶层Records数组一次批处理可携带多条记录每条记录的业务数据在kinesis子对象中database64、partitionKey、sequenceNumber、approximateArrivalTimestamp其余为事件元数据eventSource: aws:kinesis、eventIDshard 序列号、eventSourceARN流 ARN等。本地测试代码 test.py 与 test_docker.py 使用的正是这份事件结构仅invokeIdentityArn/eventSourceARN中填的是作者真实账号 ID。3. Lambda 函数源码解析核心实现在 lambda_function.py全文不足 80 行但覆盖了流式推理服务的完整要素。3.1 环境变量与模型加载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 # logged_model fruns:/{RUN_ID}/model model mlflow.pyfunc.load_model(logged_model) TEST_RUN os.getenv(TEST_RUN, False) True见 lambda_function.py三个环境变量的作用环境变量默认值作用PREDICTIONS_STREAM_NAMEride_predictions预测结果写入的输出 Kinesis 流名称RUN_ID无必填MLflow run ID用于从 S3 定位模型产物TEST_RUNFalse为True时跳过put_record只计算不写流便于本地/容器测试模型加载有两处值得注意的设计模型路径硬编码为 S3s3://mlflow-models-alexey/1/{RUN_ID}/artifacts/model中mlflow-models-alexey是课程作者实验用的 S3 桶名实际使用时需要替换为自己的桶RUN_ID则对应第 2 章Experiment Tracking中 MLflow 训练产生的 run体现了“训练产出 → S3 → 部署消费”的跨模块串联。被注释掉的runs:/{RUN_ID}/model写法则对应直接读 MLflow 模型仓库/跟踪服务器的方式。模块级加载mlflow.pyfunc.load_model在 import 阶段执行Lambda 冷启动时发生后续每次调用复用同一模型对象避免每请求重复反序列化。TEST_RUN 开关这是一个很实用的本地化设计——测试模式下不调用put_record因此本地跑test.py无需真实 AWS 凭据也能完成特征工程与预测链路验证。3.2 特征工程与预测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): pred model.predict(features) return float(pred[0])见 lambda_function.py特征与第 2、3 章训练打车时长模型时一致把上车/下车地点拼成PU_DO字符串特征供 DictVectorizer 使用外加trip_distance。返回float(pred[0])是因为输入为单条记录取预测数组第一个元素并转为原生 Python float方便json.dumps序列化。3.3 事件处理主函数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 }见 lambda_function.py处理链路是遍历event[Records]→ base64 解码 → JSON 解析 → 提取 ride 特征 → 预测 → 组装prediction_event→ 写入输出流。几个实现细节输出记录的PartitionKey选用str(ride_id)让同一行程的预测天然归入同一分区下游按ride_id消费时可以保证分区内顺序输出事件携带model与version字段——这是流式推理管道的好实践预测结果自带模型身份下游监控、审计可以区分不同模型版本产生的预测kinesis_client boto3.client(kinesis)在模块顶部创建lambda_function.py复用一个客户端返回值{predictions: predictions_events}即使在TEST_RUNTrue时也会带回全部预测供本地测试断言使用。4. 向输入流发送记录前面 2.1 节已经给出了put-record的两种形态纯文本冒烟测试和完整 ride 事件。要点--partition-key决定记录落入哪个 shard示例中固定用1--data支持内联 JSON注意外层用单引号包裹整段 JSON内部用双引号写入后可以在 Kinesis 控制台的 shard 数据中查看或等 Lambda 消费。5. 从输出流读取预测结果验证 Lambda 是否真正写入了输出流用 README 中的 CLI 片段KINESIS_STREAM_OUTPUTride_predictions SHARDshardId-000000000000 SHARD_ITERATOR$(aws kinesis \ get-shard-iterator \ --shard-id ${SHARD} \ --shard-iterator-type TRIM_HORIZON \ --stream-name ${KINESIS_STREAM_OUTPUT} \ --query ShardIterator \ ) RESULT$(aws kinesis get-records --shard-iterator $SHARD_ITERATOR) echo ${RESULT} | jq -r .Records[0].Data | base64 --decode见 streaming/README.md流程拆解get-shard-iterator用TRIM_HORIZON从 shard 最早未过期位置开始读若只想读新数据可改用LATESTget-records拉取记录注意此处输出流中的Data同样是 base64需要解码借助jq -r .Records[0].Data取第一条记录的 Data 字段管道传给base64 --decode还原 JSON——与 Lambda 内base64.b64decode(...).decode(utf-8)是同一件事的 CLI 版本。解码后应能看到形如{model: ride_duration_prediction_model, version: 123, prediction: {ride_duration: ..., ride_id: 156}}的预测事件。6. 本地测试test.py 与 test_docker.py6.1 直接调用 lambda_handlertest.py 的做法非常直接import lambda_function event { ... } # 2.3 节的完整测试事件 result lambda_function.lambda_handler(event, None) print(result)README 给出的运行方式export PREDICTIONS_STREAM_NAMEride_predictions export RUN_IDe1efc53e9bd149078b0c12aeaa6365df export TEST_RUNTrue python test.py其中RUN_ID是作者实验中的一个真实 MLflow run ID示例值。注意TEST_RUNTrue使函数跳过写流因此这一步只需要能访问 S3 上的模型需要 AWS 凭据而不需要 Kinesis 写权限。6.2 通过 HTTP 调用容器内的 Lambdatest_docker.py 则把同样的事件 POST 到本地 Lambda 模拟端点url http://localhost:8080/2015-03-31/functions/function/invocations response requests.post(url, jsonevent) print(response.json())这个 URL 正是官方 Lambda base 镜像public.ecr.aws/lambda/python:3.9内置的本地测试服务器端点2015-03-31 是 Lambda API 的版本路径容器启动后监听8080端口可以直接用任意 HTTP 客户端模拟一次真实触发。7. Docker 打包与本地容器运行7.1 Dockerfile 解析Dockerfile 全部 12 行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 ]逐行对应流式推理的打包要求基础镜像public.ecr.aws/lambda/python:3.9是 AWS 官方 Lambda base image自带provided.al2运行时入口容器启动后会执行CMD指定的 handler 并暴露 8080 本地测试端口这也是 6.2 节 URL 的来源依赖锁定pipenv install --system --deploy中的--deploy会校验Pipfile.lock与Pipfile一致否则直接失败——这对生产打包是重要的可复现性保障依赖清单Pipfile 声明boto3、mlflow、scikit-learn1.0.2Python 3.9。固定 scikit-learn 版本是为了与训练时的 pickle 产物兼容避免因版本差异导致模型加载失败。7.2 构建与运行README 中的构建与运行命令docker build -t stream-model-duration:v1 . docker run -it --rm \ -p 8080:8080 \ -e PREDICTIONS_STREAM_NAMEride_predictions \ -e RUN_IDe1efc53e9bd149078b0c12aeaa6365df \ -e TEST_RUNTrue \ -e AWS_DEFAULT_REGIONeu-west-1 \ stream-model-duration:v1测试 URLhttp://localhost:8080/2015-03-31/functions/function/invocations7.3 为容器配置 AWS 凭据README 给出两种等价方式对应 Lambda 在真实环境中“由 IAM 角色注入凭据”的两种本地模拟手段方式一注入环境变量docker run -it --rm \ -p 8080:8080 \ -e PREDICTIONS_STREAM_NAMEride_predictions \ -e RUN_IDe1efc53e9bd149078b0c12aeaa6365df \ -e TEST_RUNTrue \ -e AWS_ACCESS_KEY_ID${AWS_ACCESS_KEY_ID} \ -e AWS_SECRET_ACCESS_KEY${AWS_SECRET_ACCESS_KEY} \ -e AWS_DEFAULT_REGION${AWS_DEFAULT_REGION} \ stream-model-duration:v1方式二挂载本地.aws目录docker run -it --rm \ -p 8080:8080 \ -e PREDICTIONS_STREAM_NAMEride_predictions \ -e RUN_IDe1efc53e9bd149078b0c12aeaa6365df \ -e TEST_RUNTrue \ -v c:/Users/alexe/.aws:/root/.aws \ stream-model-duration:v1注意挂载路径示例c:/Users/alexe/.aws是 Windows 风格路径Linux/macOS 用户应替换为本机的~/.aws绝对路径容器内目标路径为/root/.aws这是因为 base image 以 root 运行。8. 发布镜像到 ECR流式管道要真正跑在 AWS 上Lambda 的镜像必须来自 ECR。README 给出了三步创建 ECR 仓库aws ecr create-repository --repository-name duration-model登录$(aws ecr get-login --no-include-email)打 tag 并推送REMOTE_URI387546586013.dkr.ecr.eu-west-1.amazonaws.com/duration-model REMOTE_TAGv1 REMOTE_IMAGE${REMOTE_URI}:${REMOTE_TAG} LOCAL_IMAGEstream-model-duration:v1 docker tag ${LOCAL_IMAGE} ${REMOTE_IMAGE} docker push ${REMOTE_IMAGE}需要强调387546586013.dkr.ecr.eu-west-1是课程作者账号与区域eu-west-1的 ECR 端点实际使用时必须替换为自己的账号 ID 和区域。镜像名沿用本地构建的stream-model-duration:v1通过 tag 映射到远程 ECR 地址后推送。推送完成后在 Lambda 控制台选择“容器镜像”包类型填入 ECR 镜像 URI即可让 Lambda 以镜像方式运行——这也正是第 6 章自动化流程所做的事。9. 进阶视角从手工流到生产级 Terraform 管道第 4 章的流式管道全靠控制台手工搭建而 06-best-practices/code 把同样的管道声明式地重建了出来可以作为“生产版”对照阅读。以 infrastructure/modules/lambda/main.tf 为例它把前面所有手工步骤变成了资源定义resource aws_lambda_function kinesis_lambda { function_name var.lambda_function_name image_uri var.image_uri # 即 ECR 中推送的镜像 package_type Image role aws_iam_role.iam_lambda.arn tracing_config { mode Active } environment { variables { PREDICTIONS_STREAM_NAME var.output_stream_name MODEL_BUCKET var.model_bucket } } timeout 180 } resource aws_lambda_function_event_source_mapping kinesis_mapping { event_source_arn var.source_stream_arn function_name aws_lambda_function.kinesis_lambda.arn starting_position LATEST ... }对照第 4 章的手工步骤可以看到一一对应关系image_uri对应 ECR 推送的镜像aws_lambda_event_source_mapping对应“Connect the function to the stream”environment.variables中注入的PREDICTIONS_STREAM_NAME与 3.1 节源码读取的环境变量同名。此外生产版还补充了两项第 4 章未涉及的能力aws_lambda_function_event_invoke_config配置了失败重试上限maximum_retry_attempts 0与最大事件年龄60 秒避免坏记录无限重试tracing_config开启 X-Ray 追踪。生产版的入口 06-best-practices/code/lambda_function.py 也体现了工程化重构handler 本体只剩一行return model_service.lambda_handler(event)其余初始化逻辑模型加载、流客户端被收敛进model.py的model.init(...)并且新增了MODEL_BUCKET环境变量以替代第 4 章硬编码的 S3 模型路径。从源码结构看这种“接口不变、内部收敛”的演进方式正是课程希望读者在第 6 章体会到的同一套流式事件契约从脚本走向可部署、可测试的服务。10. 小结与文件索引本文覆盖的流式推理管道可归纳为一条链路训练MLflow run→ 模型入 S3 → Lambda 镜像Docker ECR→ Kinesis 输入流ride_events→ 事件触发 → base64 解码 → 特征工程 → 预测 → 输出流ride_predictions→ 下游消费/监控关键文件索引均相对仓库根目录文件内容04-deployment/streaming/README.md本文主体文档流程、CLI 片段、测试事件、Docker 与 ECR 命令04-deployment/streaming/lambda_function.py核心 Lambda handler 源码04-deployment/streaming/test.py本地直接调用 handler 的测试04-deployment/streaming/test_docker.py通过 8080 端口 HTTP 调用容器内 handler 的测试04-deployment/streaming/Dockerfile / PipfileLambda base 镜像打包与依赖锁定06-best-practices/code/infrastructure/modules/lambda/main.tf生产级 Terraform镜像 Lambda 事件源映射适用前提与限制文中RUN_ID、S3 桶名mlflow-models-alexey、ECR 端点387546586013...eu-west-1均为课程作者的示例值实操时一律替换为自己的资源AWS 侧操作会产生费用这也是课程将该小节标记为可选的原因。【免费下载链接】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),仅供参考