📑 本页目录(点开跳转)
附录 D · 代码速查
📌 用法:需要什么就
Ctrl+F搜,复制粘贴。不要通读。
📦 环境
# 基础(第 1-5 节够用)
pip install numpy pandas scikit-learn
# 深度学习(第 8 节起)
pip install torch
# 向量检索(第 10 节)
pip install faiss-cpu hnswlib
# 现成的推荐库
pip install implicit lightfm recbole
# 服务化(项目三)
pip install fastapi uvicorn redis
不想装? 用 Google Colab,全部预装好了。
📊 数据加载
import pandas as pd, numpy as np, zipfile, urllib.request, io
# MovieLens 小版(1MB,秒下)
url = "https://files.grouplens.org/datasets/movielens/ml-latest-small.zip"
with urllib.request.urlopen(url) as r:
z = zipfile.ZipFile(io.BytesIO(r.read()))
ratings = pd.read_csv(z.open("ml-latest-small/ratings.csv"))
movies = pd.read_csv(z.open("ml-latest-small/movies.csv"))
⭐ 按时间切分(必须这么切,别随机切)
def time_split(df, ratio=0.8, ts_col="timestamp"):
df = df.sort_values(ts_col)
split_ts = df[ts_col].quantile(ratio)
train = df[df[ts_col] <= split_ts]
test = df[df[ts_col] > split_ts]
# 测试集只保留训练集见过的用户(冷启动单独评估)
test = test[test.userId.isin(train.userId.unique())]
return train, test
Leave-One-Out(序列推荐用)
def leave_one_out(df, user_col="userId", ts_col="timestamp"):
df = df.sort_values([user_col, ts_col])
test = df.groupby(user_col).tail(1)
valid = df.groupby(user_col).nth(-2).reset_index()
train = df.drop(test.index).drop(valid.index, errors="ignore")
return train, valid, test
构建稀疏矩阵
from scipy.sparse import csr_matrix
def build_matrix(df):
users = {u: i for i, u in enumerate(df.userId.unique())}
items = {v: i for i, v in enumerate(df.movieId.unique())}
rows = df.userId.map(users).values
cols = df.movieId.map(items).values
data = np.ones(len(df)) # 隐式反馈:有交互就是 1
mat = csr_matrix((data, (rows, cols)), shape=(len(users), len(items)))
return mat, users, items
📏 评估指标
import numpy as np
from sklearn.metrics import roc_auc_score
def recall_at_k(recs, truth, k=10):
hits = len(set(recs[:k]) & set(truth))
return hits / len(truth) if truth else 0.0
def precision_at_k(recs, truth, k=10):
return len(set(recs[:k]) & set(truth)) / k
def hit_rate_at_k(recs, truth, k=10):
return 1.0 if set(recs[:k]) & set(truth) else 0.0
def mrr(recs, truth):
for i, item in enumerate(recs):
if item in truth:
return 1.0 / (i + 1)
return 0.0
def ndcg_at_k(recs, truth, k=10):
dcg = sum(1/np.log2(i+2) for i, item in enumerate(recs[:k]) if item in truth)
idcg = sum(1/np.log2(i+2) for i in range(min(len(truth), k)))
return dcg / idcg if idcg > 0 else 0.0
# ⭐ GAUC —— 推荐系统该看的指标
def gauc(user_ids, y_true, y_pred):
user_ids, y_true, y_pred = map(np.asarray, (user_ids, y_true, y_pred))
tw, ta = 0.0, 0.0
for u in np.unique(user_ids):
m = user_ids == u
if len(np.unique(y_true[m])) < 2: # 全正或全负,跳过
continue
w = m.sum()
ta += w * roc_auc_score(y_true[m], y_pred[m])
tw += w
return ta / tw if tw else 0.0
def coverage(all_recs, n_items):
"""推荐结果覆盖了多少比例的物品库"""
return len({i for recs in all_recs for i in recs}) / n_items
# 通用评估循环
def evaluate(recommend_fn, test_dict, k=10):
res = {"recall": [], "ndcg": [], "hr": [], "mrr": []}
for user, truth in test_dict.items():
recs = recommend_fn(user, k)
if not recs: continue
res["recall"].append(recall_at_k(recs, truth, k))
res["ndcg"].append(ndcg_at_k(recs, truth, k))
res["hr"].append(hit_rate_at_k(recs, truth, k))
res["mrr"].append(mrr(recs, truth))
return {f"{m}@{k}": round(float(np.mean(v)), 4) for m, v in res.items()}
🤝 协同过滤
from collections import defaultdict
import numpy as np
class ItemCF:
"""带三个工业改进的 ItemCF:活跃用户降权 + 热门惩罚 + top-k 截断"""
def __init__(self, alpha=0.5, k=20):
self.alpha, self.k = alpha, k
def fit(self, user_items):
cooccur = defaultdict(lambda: defaultdict(float))
pop = defaultdict(int)
for user, items in user_items.items():
w = 1.0 / np.log1p(len(items)) # 活跃用户降权
for i in items:
pop[i] += 1
for j in items:
if i != j: cooccur[i][j] += w
self.sim = {}
for i, rel in cooccur.items():
s = {j: c / (pop[i]**self.alpha * pop[j]**(1-self.alpha))
for j, c in rel.items()} # 热门惩罚
self.sim[i] = dict(sorted(s.items(), key=lambda x: -x[1])[:self.k])
return self
def recommend(self, history, n=10):
scores, seen = defaultdict(float), set(history)
for i in history:
for j, s in self.sim.get(i, {}).items():
if j not in seen: scores[j] += s
return [i for i, _ in sorted(scores.items(), key=lambda x: -x[1])[:n]]
🧮 矩阵分解
用现成的(推荐)
# pip install implicit
import implicit
from scipy.sparse import csr_matrix
model = implicit.als.AlternatingLeastSquares(
factors=64, regularization=0.05, alpha=40, iterations=20)
model.fit(sparse_matrix)
ids, scores = model.recommend(userid=0, user_items=sparse_matrix[0], N=10)
# BPR 版本
model = implicit.bpr.BayesianPersonalizedRanking(factors=64, iterations=100)
model.fit(sparse_matrix)
手写 SGD 版(理解用)
见 第 5 节 的完整 MF 类。
🧠 PyTorch 模板
双塔模型
import torch, torch.nn as nn, torch.nn.functional as F
class TwoTower(nn.Module):
def __init__(self, n_users, n_items, emb=64, out=64):
super().__init__()
self.ue = nn.Embedding(n_users, emb)
self.ie = nn.Embedding(n_items, emb)
self.ut = nn.Sequential(nn.Linear(emb,128), nn.ReLU(), nn.Linear(128,out))
self.it = nn.Sequential(nn.Linear(emb,128), nn.ReLU(), nn.Linear(128,out))
def encode_user(self, u): return F.normalize(self.ut(self.ue(u)), dim=-1)
def encode_item(self, i): return F.normalize(self.it(self.ie(i)), dim=-1)
def in_batch_loss(model, users, pos_items, temp=0.05, item_freq=None):
u, v = model.encode_user(users), model.encode_item(pos_items)
logits = u @ v.T / temp
if item_freq is not None: # LogQ 校正
logits = logits - torch.log(item_freq[pos_items]).unsqueeze(0)
labels = torch.arange(len(users), device=users.device)
return F.cross_entropy(logits, labels)
通用训练循环
def train(model, loader, epochs=1, lr=1e-3, device="cpu"):
model.to(device)
opt = torch.optim.Adam(model.parameters(), lr=lr)
for ep in range(epochs):
model.train(); total = 0.0
for batch in loader:
batch = [b.to(device) for b in batch]
loss = model.loss(*batch) # 各模型自己实现 loss
opt.zero_grad(); loss.backward()
torch.nn.utils.clip_grad_norm_(model.parameters(), 5.0) # 防梯度爆炸
opt.step()
total += loss.item()
print(f"epoch {ep} loss={total/len(loader):.4f}")
return model
CTR 模型(DeepFM 简化版)
class DeepFM(nn.Module):
def __init__(self, field_dims, emb=16, hidden=[256,128]):
"""field_dims: 每个类别特征的取值个数列表,如 [10000, 500, 24, ...]"""
super().__init__()
self.offsets = np.array((0, *np.cumsum(field_dims)[:-1]))
total = sum(field_dims)
self.emb = nn.Embedding(total, emb)
self.lin = nn.Embedding(total, 1)
self.bias = nn.Parameter(torch.zeros(1))
nn.init.xavier_uniform_(self.emb.weight)
layers, in_d = [], len(field_dims) * emb
for h in hidden:
layers += [nn.Linear(in_d, h), nn.BatchNorm1d(h), nn.ReLU(), nn.Dropout(0.2)]
in_d = h
layers.append(nn.Linear(in_d, 1))
self.mlp = nn.Sequential(*layers)
def forward(self, x):
"""x: (B, num_fields) 每列是该字段的类别索引"""
x = x + x.new_tensor(self.offsets)
e = self.emb(x) # (B, F, E)
linear = self.lin(x).sum(1) + self.bias # 一阶
# FM 二阶:(Σe)² - Σ(e²) 再除 2 —— 经典化简,O(FE) 而非 O(F²E)
s = e.sum(1); fm = 0.5 * ((s*s) - (e*e).sum(1)).sum(1, keepdim=True)
deep = self.mlp(e.flatten(1))
return torch.sigmoid(linear + fm + deep).squeeze(-1)
🔍 向量检索
import faiss, numpy as np
def build_index(vectors, kind="flat"):
"""vectors: (N, D) float32"""
v = np.ascontiguousarray(vectors, dtype='float32')
faiss.normalize_L2(v) # ⭐ 内积检索前必须归一化
d = v.shape[1]
if kind == "flat": # < 10 万,精确
idx = faiss.IndexFlatIP(d)
elif kind == "hnsw": # 10 万 - 1000 万,快且准
idx = faiss.IndexHNSWFlat(d, 32, faiss.METRIC_INNER_PRODUCT)
idx.hnsw.efConstruction = 200
elif kind == "ivf": # 大规模
nlist = int(4 * np.sqrt(len(v)))
idx = faiss.IndexIVFFlat(faiss.IndexFlatIP(d), d, nlist,
faiss.METRIC_INNER_PRODUCT)
idx.train(v)
idx.add(v)
return idx
def search(index, query_vec, k=100, ef=64, nprobe=32):
q = np.ascontiguousarray(query_vec.reshape(1, -1), dtype='float32')
faiss.normalize_L2(q)
if hasattr(index, "hnsw"): index.hnsw.efSearch = ef
if hasattr(index, "nprobe"): index.nprobe = nprobe
scores, ids = index.search(q, k)
return ids[0], scores[0]
faiss.write_index(index, "items.index")
index = faiss.read_index("items.index")
🎨 重排
def mmr(scores, sim_matrix, k=10, lam=0.8):
"""lam 越大越看重相关性,越小越看重多样性"""
selected, cands = [], list(range(len(scores)))
for _ in range(min(k, len(cands))):
best, bv = None, -np.inf
for i in cands:
red = max((sim_matrix[i][j] for j in selected), default=0.0)
v = lam * scores[i] - (1 - lam) * red
if v > bv: best, bv = i, v
selected.append(best); cands.remove(best)
return selected
def scatter(items, key_fn, min_gap=3, max_per_window=2, window=5):
"""打散:同一 key(类目/作者)不要太密集"""
result, pool = [], list(items)
while pool:
placed = False
for idx, item in enumerate(pool):
k = key_fn(item)
recent = [key_fn(x) for x in result[-min_gap:]]
win = [key_fn(x) for x in result[-window:]]
if k not in recent and win.count(k) < max_per_window:
result.append(pool.pop(idx)); placed = True; break
if not placed:
result.append(pool.pop(0))
return result
🎰 探索利用
class ThompsonSampling:
def __init__(self, n_items, a=1.0, b=1.0):
self.a = np.full(n_items, a); self.b = np.full(n_items, b)
def select(self, candidates, k=10):
s = np.random.beta(self.a[candidates], self.b[candidates])
return [candidates[i] for i in np.argsort(-s)[:k]]
def update(self, item, clicked):
if clicked: self.a[item] += 1
else: self.b[item] += 1
def ucb(mean_reward, n_pulls, total_pulls, c=2.0):
if n_pulls == 0: return float('inf')
return mean_reward + c * np.sqrt(2 * np.log(total_pulls) / n_pulls)
🧪 A/B 实验
import hashlib
from scipy import stats
def assign_group(user_id, exp_id="exp_001", split=50):
"""⭐ exp_id 必须参与哈希,否则多个实验分桶重叠"""
h = int(hashlib.md5(f"{exp_id}_{user_id}".encode()).hexdigest(), 16)
return "treatment" if h % 100 < split else "control"
def ab_test(c_a, n_a, c_b, n_b, alpha=0.05):
p_a, p_b = c_a/n_a, c_b/n_b
p_pool = (c_a+c_b)/(n_a+n_b)
se = np.sqrt(p_pool*(1-p_pool)*(1/n_a + 1/n_b))
z = (p_b-p_a)/se
p_value = 2*(1-stats.norm.cdf(abs(z)))
se_d = np.sqrt(p_a*(1-p_a)/n_a + p_b*(1-p_b)/n_b)
ci = (((p_b-p_a)-1.96*se_d)/p_a, ((p_b-p_a)+1.96*se_d)/p_a)
return {"control": f"{p_a:.4%}", "treatment": f"{p_b:.4%}",
"lift": f"{(p_b-p_a)/p_a:+.2%}",
"CI95": f"[{ci[0]:+.2%}, {ci[1]:+.2%}]",
"p": round(p_value,5), "significant": p_value < alpha}
def sample_size(baseline, mde, alpha=0.05, power=0.8):
"""每组需要多少样本。mde 是相对提升,如 0.02 = 2%"""
p1, p2 = baseline, baseline*(1+mde)
za, zb = stats.norm.ppf(1-alpha/2), stats.norm.ppf(power)
pb = (p1+p2)/2
return int(np.ceil((za*np.sqrt(2*pb*(1-pb)) +
zb*np.sqrt(p1*(1-p1)+p2*(1-p2)))**2 / (p2-p1)**2))
🐛 常见报错速查
| 报错 | 原因 | 解决 |
|---|---|---|
ModuleNotFoundError: faiss |
没装 | pip install faiss-cpu |
IndexError: index out of range (Embedding) |
ID 超出 num_embeddings |
检查 ID 映射,num_embeddings = max_id + 1 |
RuntimeError: expected scalar type Long |
Embedding 输入必须是 int64 | x.long() |
CUDA out of memory |
batch 或 Embedding 太大 | 减小 batch / 减 emb_dim / torch.cuda.empty_cache() |
loss 变成 nan |
学习率太大 / log(0) | lr 除以 10;log 里加 1e-8;加梯度裁剪 |
loss 完全不降 |
lr 太小 / 数据没打乱 / 标签全一样 | 调大 lr;shuffle=True;检查标签分布 |
| Faiss 检索结果全是 -1 | 索引是空的,或 IVF 没 train | 检查 index.ntotal;IVF 必须先 train() |
| Faiss 结果不对 | 忘了 L2 归一化 | faiss.normalize_L2(v) |
roc_auc_score 报错 |
该组标签只有一类 | 跳过这组(GAUC 里已处理) |
| 推荐结果全是热门 | 没做热门惩罚 | ItemCF 加 α;召回负样本改全库随机采 |
| 离线指标高得离谱 (AUC>0.9) | ⚠️ 特征穿越 / 随机切分 | 检查是否用了未来信息;改按时间切分 |
pandas SettingWithCopyWarning |
在切片上赋值 | 用 .copy() |
🔧 性能优化速查
# ❌ 慢:循环打分
scores = [model.predict(u, i) for i in candidates]
# ✅ 快:批量打分(快 10-100 倍)
scores = model.predict_batch(u, candidates)
# ❌ 慢:循环读 Redis
feats = [redis.get(f"item:{i}") for i in items]
# ✅ 快:批量读
feats = redis.mget([f"item:{i}" for i in items])
# ❌ 慢:pandas 逐行 apply
df["x"] = df.apply(lambda r: r.a * r.b, axis=1)
# ✅ 快:向量化(快 100 倍以上)
df["x"] = df.a * df.b
# PyTorch 推理必须关梯度
with torch.no_grad():
scores = model(x)
# 或者
model.eval()
@torch.inference_mode()
def predict(x): return model(x)
🚀 FastAPI 服务模板
from fastapi import FastAPI
import hashlib, time, logging, numpy as np
app = FastAPI()
log = logging.getLogger("rec")
@app.get("/recommend")
def recommend(user_id: int, n: int = 10):
t0 = time.time()
group = assign_group(user_id)
try:
cands = list(set(vector_recall(user_id, 200)) | set(hot_recall(50)))
scored = rank_model.predict_batch(user_id, cands)
items = [i for i, _ in sorted(scored, key=lambda x: -x[1])][:n]
except Exception:
log.exception("failed, fallback")
items, group = HOT_LIST[:n], "fallback" # 🛟 绝不返回空
return {"items": items, "group": group,
"latency_ms": round((time.time()-t0)*1000, 2)}
@app.get("/health")
def health(): return {"status": "ok"}
uvicorn main:app --reload --port 8000