简介:基于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 训练好的双塔模型负责把用户和商品映射成 embedding,Milvus 负责向量召回,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、feature | to_redis.py 预热 Redis |
| products.dat | 商品 ID、类目、价格、属性 | item_id、category | rank.py 拼接特征 |
| product_vectors.txt | 商品 ID + PaddleRec 导出的 item embedding | item_id + 浮点向量 | to_milvus.py / milvus_insert.py |
| user_vector_model | 用户塔推理模型目录 | 输入 user 特征,输出 user embedding | recall.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(host=REDIS_HOST, port=REDIS_PORT, db=0, decode_responses=True) 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 = f"user:{user_id}:feat" r.set(key, json.dumps(user_feat, ensure_ascii=False))这段脚本的核心是做一个格式归一化:把 user.dat 里的行分隔字段整理成 JSON,以user:{user_id}:feat作为 key 写入 Redis。第 15 行的长度判断是用来跳过脏数据行,比如空行或者字段缺失的记录,否则后面取 fields[4] 会直接抛 IndexError。实际项目里,这个地方还会加一个异常捕获,因为用户特征里可能出现超长文本。
2.3 milvus_insert.py:把 product_vectors.txt 灌入 Milvus
to_milvus.py 和 milvus_insert.py 的分工一般是:前者负责读取 product_vectors.txt 并组装 id、向量列表,后者负责连接 Milvus、建 collection、插入和建索引。下面的片段是 milvus_insert.py 的核心逻辑:
from pymilvus import Collection, CollectionSchema, FieldSchema, DataType, connections connections.connect(host="127.0.0.1", port="19530") EMBEDDING_DIM = 64 fields = [ FieldSchema(name="item_id", dtype=DataType.INT64, is_primary=True), FieldSchema(name="embedding", dtype=DataType.FLOAT_VECTOR, dim=EMBEDDING_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.proto | UserInfoService | GetUserInfo | 按 user_id 返回用户特征 |
| item_info.proto | ItemInfoService | GetItemInfo | 按 item_id 返回商品特征 |
| recall.proto | RecallService | Recall | 返回候选 item_id 及召回分数 |
| rank.proto | RankService | Rank | 对候选列表精排打分并返回 topN |
um.py 实现 UserInfoService,从 Redis 读用户特征;cm.py 实现 ItemInfoService,查商品属性;recall.py 实现 RecallService,内部先查 Milvus;rank.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(host="127.0.0.1", port="19530") collection = Collection("item_recall") collection.load() def recall(user_embedding, top_k=50, nprobe=32): results = collection.search( data=[user_embedding], anns_field="embedding", param={"metric_type": "IP", "params": {"nprobe": nprobe}}, limit=top_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,得到每个商品的预估点击率 pCTR:
import 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, dtype="float32") 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 done"set -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 2>&1 & nohup python cm.py > logs/cm.log 2>&1 & nohup python recall.py > logs/recall.log 2>&1 & nohup python rank.py > logs/rank.log 2>&1 & nohup python as.py > logs/as.log 2>&1 & 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.py | 8003 | Redis |
| cm.py | 8004 | products.dat |
| recall.py | 8001 | Milvus、um.py |
| rank.py | 8002 | serving_server_dir |
| as.py | 8005 | recall、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_id="10001", top_k=50) resp = stub.Recall(req) for item in resp.items: print(item.item_id, round(item.score, 4))这里top_k=50表示召回 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 与模型热更新
服务能跑通只是第一步,推荐系统里更关键的是衡量线上线下不一致。这里我一般会从召回质量、检索阈值和模型更新三个角度去调。
召回质量可以用 recall@k 来验证:把测试用户真实的点击商品作为正样本集,把线上召回结果映射到同一集合,计算命中比例。下面的脚本片段可以直接跑在日志上:
def recall_at_k(reco_item_ids, click_item_ids, k=20): 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)如果 recall@20 远低于训练时的评估值,优先怀疑用户 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,避免冷启动阶段推荐结果直接为空。
本文还有配套的精品资源,点击获取