# AWS 大数据深度剖析（第 7 部分）：推荐系统基础——漏斗、双塔与 PIT

> 理解推荐系统漏斗（召回 → 粗排 → 精排 → 重排）、双塔召回架构，以及为什么 Point-in-Time（时点）正确性对训练样本至关重要。

- 作者: zhuermu
- 发布: 2025-05-16
- 网页版: https://zhuermu.com/blog/bigdata-deep-dive-part-7-recommendation-fundamentals/

---
> 数据一旦入湖，最终目标就是服务于推荐算法。本章从零讲解推荐系统：召回—排序漏斗、特征工程、模型类型，以及为什么 PIT 正确性是关键的生命线。
>
> 无需机器学习背景——只要你熟悉代码，了解“向量”和“点积”这类基本概念即可。

---

## 推荐系统解决什么问题？

当你刷 TikTok、Instagram 或 YouTube 时，每次打开 App 看到的内容都不一样——背后就有一套推荐系统在运转。

它必须回答：**“在数以亿计的候选中，应该给这个特定用户展示哪 10 个物品？并且要在 200ms 内决定。”**

挑战在于：

1. **候选太多**：平台每天收到数以千万计的新上传内容，拥有数以亿计的活跃用户
2. **延迟严苛**：用户感知到的响应必须在 200ms 以内
3. **规模化个性化**：每个用户的偏好各不相同
4. **实时反馈**：如果用户点了“不感兴趣”，下一次刷新必须立刻反映出来
5. **冷启动**：新用户没有历史行为；新内容没有互动数据

在延迟预算之内为每个用户对每个候选逐一打分排序，计算上是不可能的。**这正是我们需要一个多阶段过滤“漏斗”的原因。**

---

## 推荐漏斗：4～5 个阶段

![推荐系统漏斗](/images/blog/bigdata-deep-dive/07-recsys-funnel.svg)

### 阶段 1：候选池（10^7 量级）

整个内容库或用户图谱。这是最原始的源数据。

### 阶段 2：召回（收窄到 10^3）

快速将数以亿计的候选过滤到**几百或几千个**，供下游模型打分。

**多路召回**——并行运行 N 种不同的检索方法，然后合并结果：

| 召回通道 | 方法 | 数据来源 |
|---|---|---|
| **协同过滤**（U2U-CF / I2I-CF） | “和你相似的用户喜欢了这些” | 历史交互矩阵 |
| **双塔向量召回** | 用户向量 → 查找最近邻物品向量 | 模型训练 + 向量库 |
| **图召回**（GNN） | 在社交图谱上传播 | Neptune + Neptune ML |
| **热门召回** | 全局 / 区域热榜 | 实时统计 |
| **兴趣标签召回** | 标签匹配 | 标签倒排索引 |
| **上下文召回** | 同城 / 关注好友的内容 | 业务规则 |

每个通道检索 200～500 个物品；去重之后大约剩下 1000 个候选。得益于预计算和索引，这一步在**个位数毫秒**内完成。

### 阶段 3：粗排（收窄到 10^2）

粗排是介于召回和精排之间的中间过滤层（再缩减 10 倍）。它使用**轻量模型**进行快速打分（例如逻辑回归、浅层 MLP），以控制在延迟预算之内。

### 阶段 4：精排（收窄到 10^1）

**重型模型**（DeepFM / DIN / DCN-V2 等）对剩下的几十个候选进行精细打分。这一阶段是模型创新的主战场——几乎所有的优化精力都集中在这里。

### 阶段 5：重排

业务规则 + 多样性约束：
- 打散同类内容
- 过滤已看过的物品
- 混入广告 / 关注好友的内容
- 为冷启动新内容提供曝光保护

最终的 Top 10 会被发送到前端。

---

## 召回 vs. 排序：本质区别

很多人会把这两个阶段混为一谈。以下是关键区别：

| | 召回 | 排序 |
|---|---|---|
| 目标 | 不漏（**召回率**） | 排准（**精确率**） |
| 候选规模 | 数亿 → 数千 | 数千 → 数十 |
| 模型复杂度 | 轻量（双塔分离 + ANN） | 重型（DeepFM / Transformer） |
| 在线延迟 | 数十毫秒 | 数十毫秒 |
| 特征 | 偏全局层面的特征更多 | 偏细粒度的交叉特征更多 |
| 训练成本 | 中等 | 高 |

---

## 双塔召回模型详解

社交场景中最经典、应用最广泛的召回模型。

![双塔模型](/images/blog/bigdata-deep-dive/07-two-tower.svg)

### 模型架构

两个相互独立的神经网络（“塔”）：
- **用户塔**：消费用户特征 → 输出一个 64 维或 128 维的向量
- **物品塔**：消费物品特征 → 输出一个相同维度的向量
- 两者的**点积**或**余弦相似度** = 用户对该物品的偏好分数

```
score(user, item) = user_embedding · item_embedding
```

### 训练

样本 = (user, item, label)，其中正样本是真实的点击或关注。

**关键技巧——批内负采样（in-batch negatives）**：同一个 batch 内其他用户的正样本，作为当前用户的负样本使用（无需显式采样负样本）。

损失函数：softmax 交叉熵 / sampled softmax / BPR loss。

### 服务（关键洞见）

双塔模型之所以易于部署——因为**用户塔和物品塔是解耦的**，从而支持**离线预计算**：

```
After training:
  1. Use the item tower to precompute vectors for every item on the platform
     (tens of millions of items × 64 dimensions)
  2. Write all vectors into a vector store (OpenSearch k-NN / S3 Vectors)
  3. When a user request arrives, only run the user tower in real-time
     to compute the user vector (< 10ms)
  4. Query the vector store for the K nearest neighbors (millisecond-level ANN search)
  5. Return top-K candidate items

→ This is the essence of how two-tower enables "real-time recall":
  item vectors precomputed + ANN retrieval
```

### ANN：近似最近邻

面对数以亿计的物品向量，精确最近邻搜索代价太高。于是我们使用**近似算法**：
- HNSW（Hierarchical Navigable Small World）——当前主流
- IVF + PQ——Faiss 家族，量化可节省内存
- ScaNN（Google）

以牺牲极小精度为代价（recall@100 下降 1～2%），延迟可从秒级降到毫秒级。

> OpenSearch k-NN 支持 **Faiss / Lucene / nmslib** 后端：截至 2026 年，**Faiss 是生产环境的首选**（性能最佳 + 支持量化 + 支持 GPU）。Lucene 适合小规模、纯 JVM 的场景。nmslib 已弃用，不应再使用。

---

## 特征工程：90% 的工作都在这里

> “Garbage in, garbage out。”——特征质量决定了模型的上限。

### 特征分类

**用户侧**：
- 静态：年龄、性别、城市、注册时长
- 短期行为：最近 5 次点击、过去一小时的停留时长（**实时特征**）
- 长期偏好：常看的标签、时段规律、流量来源
- 社交：粉丝数、关注数、好友活跃度

**物品侧**：
- 静态：作者、类目、标签、创建时间
- 统计：曝光量、CTR、点赞率、完播率
- 动态：热门状态、近期评论数

**用户-物品交叉特征**：
- 用户对该作者的历史互动
- 用户对该标签的偏好分数
- 用户近期看过的相似内容

**上下文**：
- 时间（早 / 午 / 晚）、星期几、节假日
- 设备、网络（WiFi / 蜂窝）
- 地理位置

### 特征在数据仓库中如何组织

```
ads_user_features        -- User-side (daily batch update)
  user_id, age, city, last_5_click_tags, ...

ads_post_features        -- Item-side (hourly update)
  post_id, category, ctr_7d, like_rate, ...

ads_user_pair_features   -- Cross features (on demand)
  user_id, target_user_id, common_tags, ...

user_realtime_features   -- Real-time (DynamoDB / Flink maintained)
  user_id, last_click_seq[5], session_duration, ...
```

### 特征数据流

```
[ods_event] → [dwd_user_action] → [dws_user_daily] → [ads_user_features]  ← daily batch
                                                          │
                                                          ▼
                                                  Sync to DynamoDB
                                                          │
                                                          ▼
                                                  Real-time point lookup at inference

[MSK events] → Flink real-time compute last 5 clicks → DynamoDB user_realtime_features
                                                          │
                                                          ▼
                                                  Real-time point lookup at inference
```

---

## PIT 正确性（最常见的坑）

![PIT](/images/blog/bigdata-deep-dive/07-pit.svg)

### 问题

训练样本必须使用**事件发生那一刻所存在的特征取值**，而不是最新的当前取值。

否则 → **特征穿越**：模型使用了“未来”信息，离线 AUC 看起来很漂亮，但线上表现很糟糕。

### 举例

- 5 月 1 日，用户 A 点击了视频 X（一个正样本）
- 5 月 1 日，A 的兴趣标签 = “美食”
- 5 月 5 日，A 浏览视频后，算法把标签更新为“旅行”
- 5 月 10 日训练时，一个天真的 `JOIN ads_user_features ON user_id = A` → 取回的是“旅行”
- 模型学到：“旅行”用户点击美食视频 → **错误**

### 解决方案

> 警告：**纠正一个普遍存在的误解**：许多文章建议 `JOIN ads_user_features FOR TIMESTAMP AS OF s.event_ts`——**这是错误的**。在 Iceberg / SQL:2011 中，`FOR TIMESTAMP AS OF` 只接受字面常量，不接受列引用。详见 [第 02 章 2.5 节关于 PIT 的内容](./02-存储与文件格式.md#为什么推荐场景必须-iceberg-pit-正确性)。

**方案 A：每日特征快照分区（推荐，最常见）**

```sql
-- ads_user_features_daily partitioned by dt, with a full snapshot written daily
SELECT s.label, u.tag, u.age
FROM   ads_sample_follow s
JOIN   ads_user_features_daily u
       ON u.user_id = s.user_id
       AND u.dt    = s.event_dt;   -- Join on the day's snapshot
```

精度：天级。使用 Iceberg 的 `expire_snapshots` 来控制存储成本。

**方案 B：缓慢变化维（SCD Type 2，秒级精度）**

```
user_id  tag    valid_from           valid_to
A        food   2026-04-01 00:00     2026-05-05 12:00
A        travel 2026-05-05 12:00     2999-12-31 23:59
```

```sql
SELECT s.label, u.tag
FROM   ads_sample_follow s
JOIN   ads_user_features_history u
       ON s.user_id = u.user_id
       AND s.event_ts >= u.valid_from
       AND s.event_ts <  u.valid_to;
```

精度：秒级。存储更高效（仅在发生变化时新增行）。

**方案 C：SageMaker Feature Store**

一个托管的 PIT 特征存储；调用 `get_record(record_id, event_time)`，它会自动返回该时刻对应的取值。其底层封装了方案 A/B。详见第 09 章。

### Iceberg 时间旅行：超越 PIT 的真正价值

尽管它无法做逐行的 PIT 关联，时间旅行（Time Travel）在以下场景中依然至关重要：
- **数据回滚**：从误删或误更新中恢复
- **可复现训练**：固定某个 snapshot ID，几个月后仍能产出完全相同的样本
- **审计 / 调试**：对比“昨天午夜的表”与“今天午夜的表”

→ 推荐场景的完整方案是 **Iceberg + 每日快照分区 + SCD 表**。时间旅行是一种辅助能力。

### 为什么这是关键的生命线

刚入行的机器学习工程师在构造训练样本时几乎都会犯这个错误。**离线 AUC 0.85 看起来很美，但线上 CTR 纹丝不动**——根本原因就是特征穿越。这是推荐工程的“经典陷阱”。请从方案 A/B/C 中选择——**不要走上逐行时间旅行这条根本不存在的路。**

---

## 实时特征：Flink 做什么

有些特征**短暂且高频**，需要实时维护：

| 实时特征 | 说明 |
|---|---|
| 最近 5 次点击 | 供 DIN 类模型使用的行为序列 |
| 过去一小时的停留时长 | 兴趣强度信号 |
| 当前会话的动作序列 | 短期意图建模 |
| 当前网络 / 设备 / 时段 | 上下文 |

实现方式：

```
MSK events
   │
   ▼
Flink Job (keyBy user_id)
   │
   ├─ Sliding Window (5 min)
   ├─ State: maintain each user's last N clicks list (RocksDB)
   │
   ▼
DynamoDB user_realtime_features
   user_id → { last_5_clicks: [...], session_dur: 300, ... }
   │
   ▼
Recommendation inference service point lookup (milliseconds)
```

延迟：从用户点击 → 到 DynamoDB 中可见，**个位数秒级**。

---

## 模型类型概览

### 召回模型

| 模型 | 适用场景 |
|---|---|
| **协同过滤（CF）** | 经典；只要有交互矩阵就能用 |
| **双塔** | 主流；社交场景的必备 |
| **图神经网络（GNN）** | 关系数据丰富时（社交网络）效果好 |
| **序列模型（SASRec / BERT4Rec）** | 行为序列建模 |

### 排序模型

| 模型 | 适用场景 |
|---|---|
| **LightGBM / XGBoost** | 简单可靠；强特征工程能击败许多深度模型 |
| **DeepFM** | DNN + FM；兼顾特征交叉与深度 |
| **Wide & Deep** | Google 的经典架构 |
| **DIN（Deep Interest Network）** | 阿里巴巴；在行为序列上做注意力 |
| **Transformer** | 重型；效果出色但代价高昂 |
| **MMoE / PLE** | 多任务（同时优化点击 + 完播 + 关注） |

### 面向客户场景的实践建议

POC 阶段（6 个月内）：
- 召回：协同过滤 + 双塔（OpenSearch k-NN）
- 粗排：可以省略
- 排序：LightGBM（先上）→ DeepFM
- 重排：业务规则（多样性、打散）

进阶阶段：
- GNN 召回（Neptune ML）
- DIN / SIM 排序（行为序列建模）
- 多路召回融合（学习出的权重）

---

## 离线评估指标

| 指标 | 阶段 | 含义 |
|---|---|---|
| **Recall@K** | 召回 | Top-K 检索是否命中了真实点击？ |
| **Hit Rate** | 召回 | 与 Recall 类似 |
| **AUC** | 排序 | 区分能力（0.5 = 随机，1.0 = 完美） |
| **NDCG@K** | 排序 | 带权重的排序质量 |
| **GAUC** | 排序 | 按用户分别计算 AUC 再加权（更接近线上表现） |

> 警告：**离线 AUC 高不等于线上表现好。** 常见原因：PIT 错误、数据穿越、忽视曝光偏差。**A/B 测试才是真相（ground truth）。**

---

## A/B 测试

部署一个模型并不意味着立即向所有用户全量放开。你必须做 A/B 测试：

```
All users
  ├─ 50% (Control group A) → Old model
  └─ 50% (Experiment group B) → New model

Observe for N days, compare core business metrics:
  - CTR (Click-Through Rate)
  - User dwell time
  - Retention rate
  - GMV / Business conversion
```

如果新模型显著更优（统计上 p < 0.05 且业务指标提升）→ 全量放开。

A/B 测试平台通常是自研的（GrowthBook / Optimizely / 定制方案）；本指南不展开这一话题。

---

## 本章小结

| 概念 | 一句话总结 |
|---|---|
| 推荐漏斗 | 候选池 → 召回 → 粗排 → 精排 → 重排 |
| 召回 | 不漏；数亿 → 数千；多路并行 |
| 排序 | 排准；数千 → 数十；重型模型 |
| 双塔 | 主流召回模型；物品塔预计算 + ANN 检索 |
| 特征工程 | 用户 / 物品 / 交叉 / 上下文；90% 的工作 |
| PIT | 使用事件发生时刻的特征取值；善用 Iceberg 时间旅行 |
| 实时特征 | Flink 维护短期高频特征；写入 DynamoDB |

---

## 常见问题

### 什么是推荐系统漏斗？

漏斗按阶段逐层收窄候选集：召回（通过多个通道从 1000 万→1000 个物品）、粗排（通过轻量模型从 1000→100）、精排（通过 DeepFM 等深度模型从 100→50），以及重排（结合多样性规则从 50→10）。整体 P99 < 200ms。

### 机器学习训练中的 Point-in-Time（PIT，时点）正确性是什么？

PIT 指使用事件发生时刻的特征取值，而非当前的取值。缺乏 PIT 会导致特征穿越——模型在推理时本不该拿到的“未来”数据上进行训练，从而导致线上表现不佳。


---

## 参考资料

- [Wide & Deep Learning for Recommender Systems](https://arxiv.org/abs/1606.07792) — arXiv (Google)
- [DeepFM: A Factorization-Machine based Neural Network for CTR Prediction](https://arxiv.org/abs/1703.04247) — arXiv
- [Deep Learning Recommendation Model (DLRM)](https://arxiv.org/abs/1906.00091) — arXiv (Meta)
