ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

PaddleRec+Milvus电商推荐系统Linux部署实战

PaddleRec+Milvus电商推荐系统Linux部署实战 简介基于PaddleRec的深度学习电商推荐系统完整项目明确面向Linux环境部署可作为推荐系统入门实践、课程设计或毕业设计参考也可在源码基础上二次开发。压缩包共43个文件以26个Python脚本为核心搭配7个proto接口定义、启动部署Shell脚本、用户与商品数据文件、README说明文档及许可证文件包体仅43KB结构清晰便于快速下载与学习。项目覆盖PaddleRec推荐训练、召回、排序与服务化模块内置gRPC协议文件、Milvus向量召回、Redis缓存写入等配套脚本能够支撑从数据处理到模型上线的基本链路。资源已有73人学习下载代码经测试运行成功下载后按文档即可复现适合有Python基础的学生、教师或企业开发者快速实践。遇到运行问题可私聊作者获得远程教学支持也可基于现有源码扩展新的推荐策略灵活适配毕设或课设场景。1. 基于 PaddleRec 的电商推荐系统在 Linux 上到底部署了什么把这套工程解压之后第一眼会看到一堆 proto 文件、grpc 生成的 pb2 脚本以及 serving_server_dir很容易误以为只是 PaddleRec 的一个训练 demo。实际跑起来才发现它是一个完整的电商推荐系统服务端PaddleRec 训练好的双塔模型负责把用户和商品映射成 embeddingMilvus 负责向量召回Paddle Serving 负责精排打分gRPC protobuf 把召回、排序、用户特征查询拆成独立服务最后用 shell 脚本在 Linux 上一键拉起。真正花时间的地方不在模型训练而在数据格式、proto 生成顺序、服务端口和模型路径这些工程细节任何一个对不上都起不来。这套代码适合做毕设、课设也适合想在内部系统里快速落地推荐服务的开发完整跑通一遍比单独调模型收获大得多。2. 数据管线从 users.dat、products.dat 到 Redis 特征与 Milvus 向量库推荐系统在线服务的第一个瓶颈不是模型而是数据怎么被高效读取。这个工程里users.dat 和 products.dat 负责提供人的特征和商品特征product_vectors.txt 则直接保存了 PaddleRec 训练好的 item embedding。理解这三类文件的消费方式才算看懂了后面所有服务。2.1 数据文件在链路中的角色工程根目录下的几个数据文件不是摆设它们各自对接不同的下游模块。我把它们的分工整理成一张表文件内容关键字段下游消费方users.dat用户 ID、年龄、性别、城市、历史行为类目user_id、featureto_redis.py 预热 Redisproducts.dat商品 ID、类目、价格、属性item_id、categoryrank.py 拼接特征product_vectors.txt商品 ID PaddleRec 导出的 item embeddingitem_id 浮点向量to_milvus.py / milvus_insert.pyuser_vector_model用户塔推理模型目录输入 user 特征输出 user embeddingrecall.py 在线向量化users.dat 的每一行代表一个用户的静态特征和统计特征to_redis.py 会把它们转成 JSON 写入 Redis服务端拿到 user_id 后直接读 Redis避免每次请求都全量扫文件。products.dat 主要给精排阶段使用Rank 服务要把当前候选商品的特征和用户特征拼成一个稠密向量再喂给排序模型。product_vectors.txt 是已经训好的商品向量文件格式通常是第一列是商品 ID后面跟着固定维度的浮点数这一文件决定了 Milvus collection 的维度。2.2 to_redis.py用户特征离线预热到 Redis在线推荐请求的延迟预算通常在几十毫秒以内不可能让每个请求去读 users.dat 再解析。常见做法是在服务启动前把用户特征批量写入 Redis线上只做一次 key 查询。to_redis.py 干的就是这件事import redis import json REDIS_HOST 127.0.0.1 REDIS_PORT 6379 USER_FILE users.dat r redis.Redis(hostREDIS_HOST, portREDIS_PORT, db0, decode_responsesTrue) with open(USER_FILE, r) as f: for line in f: fields line.strip().split(\t) if len(fields) 5: continue user_id fields[0] user_feat { age: fields[1], gender: fields[2], city: fields[3], category_hist: fields[4], } key fuser:{user_id}:feat r.set(key, json.dumps(user_feat, ensure_asciiFalse))这段脚本的核心是做一个格式归一化把 user.dat 里的行分隔字段整理成 JSON以user:{user_id}:feat作为 key 写入 Redis。第 15 行的长度判断是用来跳过脏数据行比如空行或者字段缺失的记录否则后面取 fields[4] 会直接抛 IndexError。实际项目里这个地方还会加一个异常捕获因为用户特征里可能出现超长文本。2.3 milvus_insert.py把 product_vectors.txt 灌入 Milvusto_milvus.py 和 milvus_insert.py 的分工一般是前者负责读取 product_vectors.txt 并组装 id、向量列表后者负责连接 Milvus、建 collection、插入和建索引。下面的片段是 milvus_insert.py 的核心逻辑from pymilvus import Collection, CollectionSchema, FieldSchema, DataType, connections connections.connect(host127.0.0.1, port19530) EMBEDDING_DIM 64 fields [ FieldSchema(nameitem_id, dtypeDataType.INT64, is_primaryTrue), FieldSchema(nameembedding, dtypeDataType.FLOAT_VECTOR, dimEMBEDDING_DIM), ] schema CollectionSchema(fields, item embedding for recall) collection Collection(item_recall, schema) item_ids [] vectors [] with open(product_vectors.txt, r) as f: for line in f: parts line.strip().split() if len(parts) ! 1 EMBEDDING_DIM: continue item_ids.append(int(parts[0])) vectors.append([float(x) for x in parts[1:]]) collection.insert([item_ids, vectors]) collection.flush() index_param {index_type: IVF_FLAT, metric_type: IP, params: {nlist: 1024}} collection.create_index(embedding, index_param)这段代码里EMBEDDING_DIM必须和训练时 item embedding 的维度一致否则 insert 阶段会报 dimension mismatch。第 17 行的长度检查能提前拦截格式错误的行比如向量被截断或多了空格由于商品 ID 和向量之间可能是空格或制表符split()不传参数会把连续空白都当成分隔符这是解析文本向量最稳的写法。metric_type用 IP 是因为双塔模型做召回时常用内积近似相似度如果换 L2 距离后面 recall.py 里的 score 含义也要跟着变。nlist是 IVF 聚类的桶数量数据量小时用 1024 问题不大商品量大到千万级再往上调。2.4 维度、索引与阈值的配合这套工程里 recall 质量的上限由三件事决定向量维度、索引类型、检索时的 nprobe。维度太低模型表达能力不够不同商品的向量容易挤在一起维度太高Milvus 检索的内存和耗时都会上去。项目里 product_vectors.txt 的向量维度要跟 user_vector_model 的 user embedding 输出维度保持完全一致否则 recall.py 拿用户向量去查 item collection 时Milvus 会直接拒绝搜索。遇到这种情况优先检查训练配置里的 embedding_size而不是怀疑 Milvus 装错了。提示Milvus 1.x 和 2.x 的 Python API 差异很大。这套工程里的 Collection 写法对应 2.x如果用 1.x要先升级 pymilvus 并按 2.x 的 schema 重新建表。3. protobuf 定义接口recall.py、rank.py 背后的 gRPC 服务拆分工程里散落着 recall.proto、rank.proto、user_info.proto、item_info.proto以及对应的*_pb2.py和*_pb2_grpc.py。这些文件不是 PaddleRec 自动生成的而是把推荐服务拆成多个微服务的接口契约。先理解为什么用 protobuf gRPC再去看 recall.py 和 rank.py思路会顺很多。3.1 为什么在线服务选了 gRPC 而不是 HTTP JSON推荐系统在线链路对延迟敏感一个用户请求要经过召回、排序、特征拼接多跳JSON 序列化和反序列化的开销虽然在单次调用里不大但在高并发下会被放大。protobuf 是二进制编码体积小、解析快字段带编号前后端改动字段时不会像 JSON 那样容易静默错位。gRPC 基于 HTTP/2支持连接复用和流式传输多个请求可以共享一条 TCP 连接对长连接模型很友好。这套工程里 recall.py、rank.py、um.py 各自监听一个端口as.py 做聚合入口内部通过 gRPC stub 去调用其他服务。如果全部塞进一个 Flask 服务开发和调试简单但线上扩展、模型热更新都会受限。把召回和排序拆开后续可以单独给召回扩容也可以只对 rank 服务做灰度。3.2 四个 proto 文件划分出的服务边界proto 文件定义了服务方法、请求和响应结构。工程里几个核心服务的关系如下proto 文件服务名核心方法职责user_info.protoUserInfoServiceGetUserInfo按 user_id 返回用户特征item_info.protoItemInfoServiceGetItemInfo按 item_id 返回商品特征recall.protoRecallServiceRecall返回候选 item_id 及召回分数rank.protoRankServiceRank对候选列表精排打分并返回 topNum.py 实现 UserInfoService从 Redis 读用户特征cm.py 实现 ItemInfoService查商品属性recall.py 实现 RecallService内部先查 Milvusrank.py 实现 RankService调用排序模型得出 pCTR。as.py 则作为 API Service把客户端请求串起来。这样的分层让每个服务可以独立重启、独立部署。3.3 从 proto 生成 pb2 代码的固定步骤拿到新环境后proto 文件需要重新生成一次 Python 代码。生成顺序不能乱命令如下python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. recall.proto python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. rank.proto python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. user_info.proto python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. item_info.proto--python_out生成*_pb2.py里面是消息结构--grpc_python_out生成*_pb2_grpc.py里面是 Stub 和服务端基类。注意-I.指定 proto 的搜索路径是当前目录如果 proto 之间有互相 import比如 rank.proto 引用了 item_info.proto那么必须在同一条命令里保证所有依赖文件都在 -I 指定的路径下否则生成出来的 pb2 文件 import 路径会是错的。这个工程里源码包已经带了一份生成好的 pb2但如果从 git 拉的新环境没有安装 grpcio-tools仍然要按上面的方法重新生成。3.4 recall.py向量召回与协同过滤兜底recall.py 的核心逻辑是拿到 user embedding去 Milvus 里检索最相似的 item 向量返回候选集合。用户向量来自 user_vector_model由 user_info 服务先拼好特征再做推理。milvus_recall.py 把检索逻辑封装成了函数from pymilvus import Collection, connections connections.connect(host127.0.0.1, port19530) collection Collection(item_recall) collection.load() def recall(user_embedding, top_k50, nprobe32): results collection.search( data[user_embedding], anns_fieldembedding, param{metric_type: IP, params: {nprobe: nprobe}}, limittop_k, ) items [(str(hit.id), hit.score) for hit in results[0]] return itemscollection.load()这一步很关键Milvus 2.x 里 collection 默认不会加载到内存必须先 load 才能 search。nprobe是 IVF 索引检索时访问的聚类桶数量值越大召回越全但延迟也越高通常从 16 开始调在线 p99 延迟超了再往下减。limit是候选集大小不是最终返回给用户的条数精排 rank 阶段会从这个候选池里再截断打分。实际项目里如果用户向量检索结果为空比如新注册用户没有任何行为recall.py 会走一层协同过滤算法的兜底用热门商品或相似用户点击过的商品补齐候选。3.5 rank.py把召回候选重新排序召回阶段追求召回率排序阶段追求精准。rank.py 会把 recall 返回的候选 item 和用户特征拼接成一个完整的特征向量送入 Paddle Serving 加载的 rank_model得到每个商品的预估点击率 pCTRimport numpy as np from rank_pb2 import RankRequest, RankResponse def rank(user_emb, item_embs, model_client): features [] for item_emb in item_embs: concat_vec np.concatenate([user_emb, item_emb]).astype(float32) features.append(concat_vec) features np.array(features, dtypefloat32) scores model_client.predict(features) return scores这里的核心是特征拼接顺序必须和训练时一致训练时如果先拼 user embedding 再拼 item embedding线上推断也要保持这个顺序否则模型看到的特征分布完全错乱。model_client.predict返回的通常是一个二维数组形状是[batch_size, 1]每一行对应当前候选商品的得分。拿到 scores 后rank.py 会按分数降序排列截取 topN 返回给 as.py 聚合输出。4. Linux 部署实战prepare_server.sh、start_server.sh 与 serving_server_dir 联调这套工程在 Linux 上的启动路径很清晰先装环境再跑 prepare_server.sh 生成代码和灌数据最后 start_server.sh 拉服务。但真正操作时大多数报错都发生在环境版本、端口占用、模型加载路径这几个地方。下面按实际执行顺序拆开讲。4.1 环境准备与版本坑需要安装 Python 3.7 或 3.8、paddlepaddle、paddlerec、pymilvus、redis、grpcio、grpcio-tools。paddlepaddle 的版本要跟训练产出的模型匹配如果模型是 Paddle 2.x 导出的装 paddlepaddle 2.x 以上版本。另一个常见坑是 pymilvus 版本和 Milvus 服务端版本不一致2.x 的客户端不能连 1.x 服务端。pip install paddlepaddle pip install paddlerec pip install pymilvus redis grpcio grpcio-tools pip install paddle-serving-app paddle-serving-client paddle-serving-serverpaddle-serving 相关包在 PyPI 上的分发方式经常变装之前先确认当前环境里已经装好了对应版本的 paddlepaddle。grpcio-tools 一定要装否则 prepare_server.sh 里生成 pb2 的步骤会直接失败。装完后用python -c import redis, pymilvus, grpc; print(ok)验证导入缺哪个补哪个不要一次性装一堆不确定的依赖。4.2 prepare_server.sh生成 proto 代码并灌入数据prepare_server.sh 做的事情可以分成三段编译 proto、创建 Milvus collection 并插入商品向量、把用户特征预热到 Redis。这样一个命令就能把离线数据全部准备到位#!/bin/bash set -e echo generate grpc code... python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. recall.proto python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. rank.proto python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. user_info.proto python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. item_info.proto echo init milvus collection... python milvus_insert.py echo warm up redis... python to_redis.py echo prepare doneset -e让脚本在任一步失败时立即退出避免后面服务起来后才发现数据没灌进去。这里要特别注意执行顺序milvus_insert.py 依赖 to_milvus.py 生成的向量列表而 to_milvus.py 又依赖 product_vectors.txt 的文件格式。如果商品向量文件还没生成Milvus 里就是空集合recall 服务查不到任何候选。redis 预热放在最后因为它只影响用户特征读取不影响 proto 生成。4.3 start_server.sh按依赖顺序拉起服务服务启动顺序是先起用户特征服务和商品特征服务再起召回和排序最后起 as.py 聚合入口。start_server.sh 里常见的写法如下#!/bin/bash mkdir -p logs nohup python um.py logs/um.log 21 nohup python cm.py logs/cm.log 21 nohup python recall.py logs/recall.log 21 nohup python rank.py logs/rank.log 21 nohup python as.py logs/as.log 21 sleep 3 python client.py每个服务监听不同端口um.py 和 cm.py 相对独立recall.py 依赖 Milvus 和 um.py 提供的用户特征rank.py 依赖模型文件加载。最后 sleep 3 秒是为了等 as.py 起来后client.py 请求不会直接 connection refused。服务端口划分如下服务默认端口依赖um.py8003Rediscm.py8004products.datrecall.py8001Milvus、um.pyrank.py8002serving_server_diras.py8005recall、rank、um、cm我把 as.py 的聚合逻辑理解为对外 API 层客户端只连它一个端口内部再并发调用召回和排序。这样如果要把 gRPC 换成 HTTP只需要改 as.py下游不用动。4.4 client.py 验证整条链路client.py 是测试工具它会把一个真实请求打到 as.py再打印返回的商品列表。如果用 gRPC stub 直接连 recall 服务请求写法类似下面这样import grpc import recall_pb2 import recall_pb2_grpc channel grpc.insecure_channel(127.0.0.1:8001) stub recall_pb2_grpc.RecallServiceStub(channel) req recall_pb2.RecallRequest(user_id10001, top_k50) resp stub.Recall(req) for item in resp.items: print(item.item_id, round(item.score, 4))这里top_k50表示召回 50 个候选最终用户看到的推荐列表通常小于这个数字因为 rank 阶段会结合 pCTR 再截断到 10 或 20。如果直接连 as.py 的端口它会内部帮你串起召回和排序返回的是 final topN。调试时建议先用 client.py 直连 recall 端口确认召回有数据再切到 as.py这样能快速定位问题在召回还是排序阶段。4.5 服务起不来时的排查套路最常见的问题是端口被占用。不同的服务监听了多个端口第二次跑 start_server.sh 时经常出现Address already in use用 lsof 查端口占用lsof -i:8001 lsof -i:8002如果端口被残留进程占用先 kill 再重启不要硬改端口号因为 as.py 和 client.py 里的端口配置是连在一起的。其次是 proto import 报错比如ModuleNotFoundError: No module named recall_pb2说明当前工作目录不在 Python 的搜索路径里运行 src 下的服务脚本时要cd到工程根目录或在脚本头部加sys.path.append(os.path.dirname(__file__))。Milvus 相关报错集中在 collection not found 和 dimension mismatch 两类前者说明 prepare 阶段没跑成功后者说明 product_vectors.txt 的维度和 milvus_insert.py 里的 EMBEDDING_DIM 不一致。start_server.sh 里的 logs 目录下每个服务的日志都是独立文件出问题先grep -i error logs/*.log。5. 线上验证与再调优向量阈值、nprobe 与模型热更新服务能跑通只是第一步推荐系统里更关键的是衡量线上线下不一致。这里我一般会从召回质量、检索阈值和模型更新三个角度去调。召回质量可以用 recallk 来验证把测试用户真实的点击商品作为正样本集把线上召回结果映射到同一集合计算命中比例。下面的脚本片段可以直接跑在日志上def recall_at_k(reco_item_ids, click_item_ids, k20): reco_set set(reco_item_ids[:k]) click_set set(click_item_ids) if not click_set: return 0.0 return len(reco_set click_set) / len(click_set)如果 recall20 远低于训练时的评估值优先怀疑用户 embedding 和 item embedding 的分布不一致检查 user_vector_model 是否用的是同一套特征处理逻辑。Milvus 检索参数也值得调nprobe从 16 调到 64 通常能明显提升召回率但如果召回集合变大后精排来不及处理就要通过设置 score 阈值过滤低质量候选。比如 IP 内积分数低于某个值的 item 直接丢弃这个阈值我习惯按线上日志的分布动态调先记 0.5 起步再往下探。模型热更新不需要把整套服务重启。rank_model 更新时把新模型文件替换到 serving_server_dir 下然后单独重启 rank.py 进程即可recall 侧的 user_vector_model 更新后要确认线上推理输出的向量维度和 Milvus collection 一致否则查询会失败。新商品没有向量时最直接的做法是用 products.dat 里同类目的热门商品补位等离线任务重新产出一批 embedding 后再洗入 Milvus避免冷启动阶段推荐结果直接为空。本文还有配套的精品资源点击获取
返回列表