PyTorch实战:手把手教你构建用户-商品Embedding推荐系统(附完整代码)
从零构建用户-商品Embedding推荐系统:PyTorch实战与深度调优指南
推荐系统早已不是新鲜概念,但如何让机器真正“理解”用户的偏好,并精准匹配到合适的商品,依然是许多开发者面临的挑战。如果你已经掌握了PyTorch的基础,却苦于不知如何将理论知识落地为一个高效、可扩展的推荐模型,那么这篇文章正是为你准备的。我们将绕过那些泛泛而谈的理论,直接切入核心,手把手带你构建一个基于Embedding的用户-商品推荐系统。整个过程不仅包含可运行的完整代码,更会深入探讨数据处理的陷阱、模型设计的权衡、训练技巧的细节,以及如何将模型部署到实际场景中进行评估和迭代。无论你是希望为自己的项目增加推荐功能,还是想深入理解现代推荐系统的核心机制,这里都有你需要的实战经验。
1. 项目蓝图:理解我们要构建什么
在动手写第一行代码之前,我们必须清晰地勾勒出整个项目的轮廓。一个典型的基于Embedding的推荐系统,其核心思想异常直观:将用户和商品都映射到同一个低维的向量空间中。在这个空间里,向量之间的距离或方向就代表了用户兴趣与商品特性之间的匹配程度。
想象一下,我们把所有用户和商品都放在一个巨大的“兴趣宇宙”里。喜欢科幻电影和电子产品的用户,他们的向量会聚集在宇宙的某个区域;而与之相关的商品,比如《沙丘》蓝光碟和最新游戏显卡,其向量也应该出现在附近。推荐,本质上就是在这个宇宙中,为用户寻找最近邻的商品。
我们的项目将遵循一个清晰的流水线:
- 数据层:处理原始的用户行为日志(点击、购买、评分)。
- 模型层:构建PyTorch模型,学习用户和商品的Embedding向量。
- 训练层:设计损失函数和采样策略,高效地训练Embedding。
- 评估与应用层:评估模型效果,并实现简单的推荐查询。
与许多教程不同,我们将特别关注负采样策略的工程实现和离线评估指标的陷阱,这些往往是项目成败的关键。
2. 数据工程:从原始日志到模型可用的样本
数据决定了模型的天花板。我们假设你拥有最典型的隐式反馈数据:一系列(user_id, item_id, timestamp)三元组,代表用户对商品产生了某种正向行为(如点击)。没有负样本(即用户明确不喜欢的记录),这正是隐式反馈推荐的特点和难点。
2.1 数据预处理与关键洞察
首先,我们需要将原始数据加载并转化为结构化的DataFrame。
import pandas as pd
import numpy as np
from collections import defaultdict
import warnings
warnings.filterwarnings('ignore')
# 模拟生成一些数据
np.random.seed(42)
num_users = 1000
num_items = 5000
num_interactions = 50000
# 生成随机交互数据
user_ids = np.random.randint(0, num_users, size=num_interactions)
item_ids = np.random.randint(0, num_items, size=num_interactions)
# 假设交互强度(如观看时长归一化到0-1),可作为训练权重
strength = np.random.rand(num_interactions)
df_interactions = pd.DataFrame({
'user_id': user_ids,
'item_id': item_ids,
'strength': strength,
'timestamp': np.arange(num_interactions) # 简单递增的时间戳
})
print(f"原始交互数据概览:\n{df_interactions.head()}")
print(f"总交互数:{len(df_interactions)}")
接下来是至关重要的分析步骤:检查数据的稀疏性和分布。一个高度稀疏的数据集(即大多数用户只与极少数商品交互)需要不同的处理策略。
# 计算数据稀疏度
sparsity = 1 - len(df_interactions) / (num_users * num_items)
print(f"数据稀疏度:{sparsity:.4%}")
# 分析用户和商品交互的分布
user_interaction_counts = df_interactions['user_id'].value_counts()
item_interaction_counts = df_interactions['item_id'].value_counts()
print(f"\n用户平均交互数:{user_interaction_counts.mean():.2f}")
print(f"商品平均被交互数:{item_interaction_counts.mean():.2f}")
print(f"最活跃用户交互数:{user_interaction_counts.max()}")
print(f"最热门商品被交互数:{item_interaction_counts.max()}")
注意:长尾分布(即少数用户/商品占据了大部分交互)在推荐系统中非常普遍。对于交互次数极少的“冷启动”用户或商品,其Embedding难以训练,后续可能需要特殊处理,如使用全局平均向量进行初始化。
2.2 构建训练样本:正样本与负采样的艺术
对于每个正向交互(u, i),我们都需要为其生成负样本。一种简单但有效的策略是:为每个正样本,随机采样若干个该用户未交互过的商品作为负样本。这里的关键在于“随机”并非完全均匀,有时会采用“流行度偏置”采样(即更热门的商品有更高概率被选为负样本),以帮助模型更好地区分流行商品和用户真实兴趣。
我们先构建一个高效的数据结构,用于快速查询用户交互过的商品集合。
# 构建用户-商品交互字典
user_interacted_items = defaultdict(set)
for _, row in df_interactions.iterrows():
user_interacted_items[row['user_id']].add(row['item_id'])
# 构建商品全集列表,用于负采样
all_item_ids = set(range(num_items))
现在,实现一个生成训练批次的函数。我们将采用经典的BPR(Bayesian Personalized Ranking) 损失所需的样本格式:对于每个用户,我们有一个正商品和一个负商品。
def generate_bpr_batch(df, user_interacted_dict, all_items_set, batch_size=512, neg_per_pos=1):
"""
生成BPR训练所需的批次数据。
返回:user_ids, positive_item_ids, negative_item_ids
"""
users = []
pos_items = []
neg_items = []
# 随机打乱数据
df_shuffled = df.sample(frac=1).reset_index(drop=True)
for idx, row in df_shuffled.iterrows():
u = row['user_id']
i = row['item_id']
# 获取该用户未交互的商品池
interacted = user_interacted_dict[u]
non_interacted = list(all_items_set - interacted)
if len(non_interacted) == 0:
continue # 如果用户交互了所有商品(几乎不可能),则跳过
# 为每个正样本采样neg_per_pos个负样本
for _ in range(neg_per_pos):
j = np.random.choice(non_interacted)
users.append(u)
pos_items.append(i)
neg_items.append(j)
# 如果收集的样本数达到batch_size,则yield一个批次
if len(users) >= batch_size:
yield (np.array(users[:batch_size]),
np.array(pos_items[:batch_size]),
np.array(neg_items[:batch_size]))
users = users[batch_size:]
pos_items = pos_items[batch_size:]
neg_items = neg_items[batch_size:]
# 返回最后不足一个batch的剩余数据
if len(users) > 0:
yield (np.array(users), np.array(pos_items), np.array(neg_items))
3. 模型架构:用PyTorch定义Embedding与交互逻辑
我们将构建一个灵活的双塔模型(Two-Tower Model)。用户塔和商品塔独立计算Embedding,最后通过一个交互函数(如点积)计算匹配分数。这种结构便于线上服务时缓存商品向量,实现毫秒级推荐。
3.1 基础双塔模型实现
import torch
import torch.nn as nn
import torch.nn.functional as F
class TwoTowerModel(nn.Module):
"""
基础双塔模型。
用户侧和商品侧共享相同的Embedding维度,但参数独立。
"""
def __init__(self, num_users, num_items, embedding_dim=64, hidden_dims=[128, 64]):
super(TwoTowerModel, self).__init__()
self.embedding_dim = embedding_dim
# 用户塔
self.user_embedding = nn.Embedding(num_users, embedding_dim)
user_layers = []
input_dim = embedding_dim
for h_dim in hidden_dims:
user_layers.append(nn.Linear(input_dim, h_dim))
user_layers.append(nn.ReLU())
user_layers.append(nn.BatchNorm1d(h_dim))
input_dim = h_dim
self.user_tower = nn.Sequential(*user_layers)
# 商品塔(结构可以与用户塔对称或不同)
self.item_embedding = nn.Embedding(num_items, embedding_dim)
item_layers = []
input_dim = embedding_dim
for h_dim in hidden_dims:
item_layers.append(nn.Linear(input_dim, h_dim))
item_layers.append(nn.ReLU())
item_layers.append(nn.BatchNorm1d(h_dim))
input_dim = h_dim
self.item_tower = nn.Sequential(*item_layers)
# 初始化Embedding权重
self._init_weights()
def _init_weights(self):
""" Xavier均匀初始化,有助于训练稳定性 """
nn.init.xavier_uniform_(self.user_embedding.weight)
nn.init.xavier_uniform_(self.item_embedding.weight)
for layer in self.user_tower:
if isinstance(layer, nn.Linear):
nn.init.xavier_uniform_(layer.weight)
for layer in self.item_tower:
if isinstance(layer, nn.Linear):
nn.init.xavier_uniform_(layer.weight)
def forward(self, user_ids, item_ids):
"""
前向传播。
返回:用户向量,商品向量,匹配分数(点积)
"""
user_emb = self.user_embedding(user_ids)
item_emb = self.item_embedding(item_ids)
user_vec = self.user_tower(user_emb)
item_vec = self.item_tower(item_emb)
# 对向量进行L2归一化,使点积等于余弦相似度
user_vec = F.normalize(user_vec, p=2, dim=1)
item_vec = F.normalize(item_vec, p=2, dim=1)
# 计算匹配分数(余弦相似度)
score = (user_vec * item_vec).sum(dim=1)
return user_vec, item_vec, score
def get_user_embedding(self, user_ids):
""" 单独获取用户向量,用于批量预计算或缓存 """
user_emb = self.user_embedding(user_ids)
user_vec = self.user_tower(user_emb)
return F.normalize(user_vec, p=2, dim=1)
def get_item_embedding(self, item_ids):
""" 单独获取商品向量,用于批量预计算或缓存 """
item_emb = self.item_embedding(item_ids)
item_vec = self.item_tower(item_emb)
return F.normalize(item_vec, p=2, dim=1)
3.2 损失函数的选择:BPR Loss详解
对于隐式反馈,BPR损失是一个经典且强大的选择。它的核心思想是最大化正样本商品与负样本商品对于同一用户的排序分差。
class BPRLoss(nn.Module):
"""
Bayesian Personalized Ranking Loss.
Loss = -log(sigmoid(pos_score - neg_score)) 的均值。
"""
def __init__(self):
super(BPRLoss, self).__init__()
def forward(self, pos_scores, neg_scores):
"""
pos_scores: 正样本的匹配分数 [batch_size]
neg_scores: 负样本的匹配分数 [batch_size]
"""
# 计算差值
diff = pos_scores - neg_scores
# 应用sigmoid并取负对数
loss = -torch.log(torch.sigmoid(diff) + 1e-8).mean()
return loss
为什么BPR有效?因为它直接优化了排序指标(AUC),假设用户更喜欢他交互过的商品胜过未交互的,而不需要关心具体分数的大小。在实际训练中,我发现对损失函数加上L2正则化能有效防止过拟合,特别是当用户或商品交互数据很少时。
4. 训练流程:工程实现与性能调优
现在我们将所有组件组装起来,形成一个完整的训练循环。这里会包含学习率调度、模型检查点保存以及简单的训练监控。
4.1 训练循环实现
from torch.utils.data import DataLoader, TensorDataset
import time
def train_epoch(model, data_generator, optimizer, criterion, device):
""" 训练一个epoch """
model.train()
total_loss = 0.0
num_batches = 0
for batch_users, batch_pos, batch_neg in data_generator:
# 将数据移动到设备
batch_users = torch.LongTensor(batch_users).to(device)
batch_pos = torch.LongTensor(batch_pos).to(device)
batch_neg = torch.LongTensor(batch_neg).to(device)
# 清零梯度
optimizer.zero_grad()
# 前向传播:分别计算正负样本的分数
_, _, pos_scores = model(batch_users, batch_pos)
_, _, neg_scores = model(batch_users, batch_neg)
# 计算损失
loss = criterion(pos_scores, neg_scores)
# 反向传播与优化
loss.backward()
# 可选:梯度裁剪,防止梯度爆炸
torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=5.0)
optimizer.step()
total_loss += loss.item()
num_batches += 1
return total_loss / max(num_batches, 1)
def prepare_dataloader(df, user_interacted_dict, all_items_set, batch_size=1024, neg_per_pos=4):
""" 将生成器转换为可迭代的数据加载器 """
def batch_generator():
yield from generate_bpr_batch(df, user_interacted_dict, all_items_set, batch_size, neg_per_pos)
return batch_generator()
# 配置训练参数
embedding_dim = 64
hidden_dims = [128, 64]
learning_rate = 0.001
num_epochs = 20
batch_size = 1024
neg_per_pos = 4
# 初始化模型、损失函数、优化器
device = torch.device('cuda' if torch.cuda.is_available() else 'cpu')
print(f"使用设备:{device}")
model = TwoTowerModel(num_users, num_items, embedding_dim, hidden_dims).to(device)
criterion = BPRLoss()
optimizer = torch.optim.Adam(model.parameters(), lr=learning_rate, weight_decay=1e-5) # weight_decay即L2正则化
scheduler = torch.optim.lr_scheduler.ReduceLROnPlateau(optimizer, mode='min', patience=2, factor=0.5, verbose=True)
# 训练循环
train_loss_history = []
print("开始训练...")
for epoch in range(num_epochs):
start_time = time.time()
# 准备当前epoch的数据生成器
train_data_gen = prepare_dataloader(df_interactions, user_interacted_items, all_item_ids, batch_size, neg_per_pos)
# 训练一个epoch
epoch_loss = train_epoch(model, train_data_gen, optimizer, criterion, device)
train_loss_history.append(epoch_loss)
# 调整学习率
scheduler.step(epoch_loss)
epoch_time = time.time() - start_time
print(f"Epoch [{epoch+1:03d}/{num_epochs}] | Loss: {epoch_loss:.6f} | Time: {epoch_time:.2f}s")
# 每隔几个epoch保存一次模型检查点
if (epoch + 1) % 5 == 0:
checkpoint_path = f'model_checkpoint_epoch_{epoch+1}.pth'
torch.save({
'epoch': epoch,
'model_state_dict': model.state_dict(),
'optimizer_state_dict': optimizer.state_dict(),
'loss': epoch_loss,
}, checkpoint_path)
print(f"模型已保存至:{checkpoint_path}")
4.2 关键训练技巧与陷阱
在训练过程中,有几个细节会显著影响最终效果:
- Embedding维度:并非越大越好。过大的维度容易在小数据集上过拟合。通常从32或64开始尝试,根据验证集效果调整。
- 负采样比例:
neg_per_pos是一个重要超参。比例太低,模型区分能力不足;比例太高,训练不稳定。一般设置在3到10之间。 - 批次大小:较大的批次通常能提供更稳定的梯度估计,但会占用更多显存。如果使用点积相似度,极大批次可能需要考虑跨批次负采样来增加每个正样本看到的负样本数。
- 向量归一化:在forward函数中对最终向量进行L2归一化,将点积转化为余弦相似度,有助于训练稳定,且分数范围固定在[-1, 1]。
提示:在真实场景中,建议划分出一部分时间上最新的交互作为验证集,用于早停(Early Stopping)和超参调优,防止模型过拟合到历史行为模式。
5. 评估与推荐:离线指标与线上服务模拟
模型训练完成后,我们需要评估其推荐质量,并实现推荐逻辑。
5.1 离线评估:召回率与NDCG
我们采用留一法(Hold-One-Out)进行评估:对于每个用户,将其最后一次交互作为测试正样本,其余用于训练(或直接使用已训练模型)。然后为该用户推荐K个商品,计算测试商品是否出现在推荐列表中。
from sklearn.metrics import ndcg_score
import numpy as np
def evaluate_model(model, df_test, df_full, user_interacted_dict, all_items_set, k=10, device='cpu'):
"""
评估模型性能。
df_test: 包含每个用户最后一次交互的DataFrame
df_full: 全量交互数据,用于排除已交互商品
"""
model.eval()
recalls = []
ndcgs = []
# 预计算所有商品向量(在实际中可缓存)
all_item_ids_tensor = torch.LongTensor(list(range(num_items))).to(device)
with torch.no_grad():
all_item_vecs = model.get_item_embedding(all_item_ids_tensor).cpu().numpy()
for _, row in df_test.iterrows():
u = row['user_id']
test_item = row['item_id']
# 获取该用户训练阶段已交互的商品(排除测试商品)
train_items = user_interacted_dict[u] - {test_item}
# 获取用户向量
u_tensor = torch.LongTensor([u]).to(device)
with torch.no_grad():
user_vec = model.get_user_embedding(u_tensor).cpu().numpy()
# 计算用户与所有商品的分数(余弦相似度)
scores = np.dot(all_item_vecs, user_vec.T).flatten()
# 将训练集中已交互的商品分数设为极低,确保不会被推荐
scores[list(train_items)] = -np.inf
# 获取top-K推荐商品的索引
top_k_indices = np.argsort(scores)[-k:][::-1]
# 计算Recall@K
recall = 1.0 if test_item in top_k_indices else 0.0
recalls.append(recall)
# 计算NDCG@K:构建相关性列表,测试商品相关度为1,其他为0
relevance = np.zeros(num_items)
relevance[test_item] = 1
# 注意:sklearn的ndcg_score期望二维输入,且从高到低排序
dcg_idcg = ndcg_score([relevance], [scores], k=k)
ndcgs.append(dcg_idcg)
mean_recall = np.mean(recalls)
mean_ndcg = np.mean(ndcgs)
return mean_recall, mean_ndcg
# 模拟构建测试集(每个用户取最后一次交互)
df_test = df_interactions.sort_values('timestamp').groupby('user_id').tail(1)
print(f"测试集大小:{len(df_test)}")
# 进行评估
recall_k, ndcg_k = evaluate_model(model, df_test, df_interactions, user_interacted_items, all_item_ids, k=10, device=device)
print(f"评估结果 (K=10) -> Recall: {recall_k:.4f}, NDCG: {ndcg_k:.4f}")
5.2 实现推荐函数
最后,我们实现一个简单的推荐函数,展示如何为指定用户生成个性化推荐列表。
def recommend_for_user(model, user_id, user_interacted_set, all_item_vecs, k=20, device='cpu'):
"""
为单个用户生成Top-K推荐。
all_item_vecs: 预计算好的所有商品向量矩阵 [num_items, embedding_dim]
"""
model.eval()
# 获取用户向量
u_tensor = torch.LongTensor([user_id]).to(device)
with torch.no_grad():
user_vec = model.get_user_embedding(u_tensor).cpu().numpy()
# 计算分数
scores = np.dot(all_item_vecs, user_vec.T).flatten()
# 屏蔽已交互商品
scores[list(user_interacted_set)] = -np.inf
# 返回Top-K商品ID及其分数
top_k_indices = np.argsort(scores)[-k:][::-1]
top_k_scores = scores[top_k_indices]
return list(zip(top_k_indices.tolist(), top_k_scores.tolist()))
# 示例:为用户0推荐商品
# 首先需要预计算商品向量(在实际服务中,此计算可离线进行并缓存)
all_item_ids_tensor = torch.LongTensor(list(range(num_items))).to(device)
with torch.no_grad():
cached_item_vecs = model.get_item_embedding(all_item_ids_tensor).cpu().numpy()
user_id_to_test = 0
recommendations = recommend_for_user(model, user_id_to_test, user_interacted_items[user_id_to_test], cached_item_vecs, k=5, device=device)
print(f"\n为用户 {user_id_to_test} 的Top-5推荐:")
for item_id, score in recommendations:
print(f" 商品ID: {item_id:4d} | 预测分数: {score:.4f}")
6. 进阶优化与生产化考量
至此,一个基础但完整的推荐系统已经构建完成。然而,要将其应用于真实生产环境,还需要考虑更多因素。
模型层面:
- 引入侧信息:当前的模型只使用了ID信息。可以扩展模型,融入用户画像(年龄、性别)和商品属性(类别、价格)的Embedding,通过多层感知机(MLP)与ID Embedding融合,提升冷启动效果。
- 序列建模:用户兴趣是动态变化的。可以考虑使用GRU或Transformer编码用户最近的交互序列,捕捉短期兴趣。
- 多目标优化:除了点击率,还可以同时优化停留时长、购买转化率等,构建多任务学习模型。
工程层面:
- 近似最近邻搜索:当商品数量达到百万甚至千万级时,遍历计算所有商品分数是不现实的。需要集成FAISS或Annoy等近似最近邻库,实现毫秒级检索。
- 模型部署与更新:如何将PyTorch模型部署为API服务?如何设计在线学习流程,使模型能随着新数据流入而持续更新?这通常需要结合TorchServe、Redis(缓存向量)和流处理平台(如Kafka)。
- A/B测试框架:任何推荐算法的改进,最终都需要通过线上A/B测试验证其业务价值。需要设计一套完整的指标埋点、实验分组和统计分析流程。
构建推荐系统是一个持续迭代的过程。从基础的Embedding双塔模型出发,理解数据、模型、训练、评估的每一个环节,再根据实际业务需求和数据特点,逐步引入更复杂的模块和优化策略,才是稳健的技术演进之路。
更多推荐
所有评论(0)