推荐系统大数据架构
概述
推荐系统是大数据 + 机器学习的典型落地场景。完整链路:特征 → 召回 → 排序 → 重排,且要离线算 + 实时算结合。本文从大数据视角讲透推荐架构,并给一个完整链路案例。
一、推荐系统总体架构
1.1 推荐链路
用户请求 → 特征 → 召回(候选集)
→ 粗排 → 精排 → 重排 → 返回
每层目的不同:
召回:从海量物品中找候选
粗排:快速过滤
精排:精细打分
重排:业务规则/多样性| 阶段 | 候选量 | 要求 |
|---|---|---|
| 召回 | 全量 → 千级 | 快、全 |
| 粗排 | 千级 → 百级 | 轻量模型 |
| 精排 | 百级 → 几十 | 强模型 |
| 重排 | 几十 → 结果 | 规则/多样性 |
1.2 大数据的作用
大数据支撑:
海量行为数据(训练基础)
大规模特征(模型输入)
实时数据(新鲜特征)
离线训练(模型更新)
离线与实时结合:
离线:模型训练、批量特征、全量打分
实时:实时特征、实时召回、秒级更新二、数据层
2.1 数据来源
| 数据 | 说明 | 采集 |
|---|---|---|
| 行为日志 | 点击/曝光/购买 | 埋点 + Kafka |
| 用户画像 | 属性/标签 | 数仓加工 |
| 物品信息 | 类目/属性 | 业务库 CDC |
| 上下文 | 时间/渠道 | 实时 |
2.2 数据处理
流程:
埋点 → Kafka → 实时(Flink 清洗/特征)
→ 数仓(离线聚合/画像)
分层:
ODS:原始行为
DWD:清洗明细
DWS:行为聚合(曝光/点击/购买)
ADS:特征与标签2.3 特征体系
特征分类:
用户特征(画像/历史偏好)
物品特征(类目/热度/内容)
行为特征(近 7 天交互)
交叉特征(用户×物品)
存储:
离线特征(Feature Store)
实时特征(Redis/HBase)三、召回层
3.1 召回方法
| 方法 | 原理 | 特点 |
|---|---|---|
| 协同过滤 | 用户/物品相似 | 经典 |
| 向量召回 | Embedding 相似(ANN) | 主流 |
| 规则召回 | 热门/新品/类目 | 兜底 |
| 图召回 | 知识图谱 | 冷启 |
向量召回:
用户/物品 → Embedding 向量
相似度(内积/余弦)
近邻检索(Faiss/ANN)
流程:
离线算向量 → 向量索引
在线:用户向量 → 检索 TOP-K3.2 召回架构
召回编排:
多个召回源并行:
协同过滤候选
向量召回候选
热门/新品候选
→ 合并去重 → 粗排
大数据支撑:
离线训练 Embedding
构建向量索引
全量物品向量更新四、排序层
4.1 粗排
粗排目的:
从千级候选快速筛到百级
方案:
轻量模型(LR/双塔)
规则打分
要求:
快(低算力)
别把好的滤掉4.2 精排
精排模型:
特征(用户/物品/行为/交叉)
→ 模型打分(点击率/转化率)
常用模型:
GBDT/XGBoost(传统)
DeepFM/Wide&Deep(深度学习)
多目标(CTR/CVR 联合)
特点:
特征工程重要
训练数据大(亿级样本)4.3 排序的实时性
离线模型 + 实时特征:
模型:日级/小时级更新
特征:实时行为秒级生效
实时排序:
实时特征接入模型
模型服务低延迟打分
排序结果反馈回流五、重排层
5.1 重排目标
重排考虑:
多样性(类目不单一)
打散(同作者/同店铺限流)
业务规则(广告位/库存)
惊喜度(探索新物品)
手段:
规则重排
强化学习(进阶)
多样性与相关性权衡5.2 重排实现
流程:
精排结果 → 业务规则过滤
→ 多样性打散
→ 位置约束
→ 最终列表
常用算法:
最大边际相关(MMR)
贪心打散六、离线与实时计算
6.1 离线计算
离线任务:
训练样本生成(行为数据聚合)
模型训练(日级/小时级)
批量特征更新
物品向量/索引构建
全量离线打分(预生成)
平台:
调度(Airflow/DolphinScheduler)
Spark 数据处理
训练平台6.2 实时计算
实时任务:
实时行为特征(Flink 窗口聚合)
实时召回(实时协同/热度)
实时画像更新
效果实时监控
链路:
埋点 → Kafka → Flink → 特征库
→ 在线服务读取6.3 在线服务
在线链路:
请求 → 特征服务 → 召回 → 排序
→ 重排 → 返回
性能:
P99 延迟目标(如 200ms)
缓存(特征/结果)
异步并行(召回/特征并发)七、全链路案例
7.1 电商推荐架构案例
整体链路:
用户打开 APP
→ 请求推荐服务
→ 取实时特征(Redis)
→ 多路召回(向量/热门/协同)
→ 粗排(双塔打分)
→ 精排(DeepFM CTR 预估)
→ 重排(多样性 + 广告)
→ 返回商品列表
离线支撑:
日级训练样本(Spark 聚合)
模型训练(日级)
物品向量更新(小时级)
全量召回索引(Faiss)
实时支撑:
点击事件 → Kafka → Flink
→ 实时特征更新(Redis)
→ 实时热度召回7.2 组件清单
| 环节 | 组件 |
|---|---|
| 采集 | 埋点 + Kafka |
| 实时计算 | Flink |
| 离线处理 | Spark |
| 特征库 | Redis + Feature Store |
| 向量检索 | Faiss |
| 模型服务 | TensorRT/Triton |
| 调度 | Airflow/DolphinScheduler |
| 监控 | Prometheus + Grafana |
7.3 关键指标
效果指标:
点击率(CTR)
转化率(CVR)
人均点击/成交
冷启动效果
系统指标:
延迟 P99
召回覆盖率
特征新鲜度八、常见问题
8.1 冷启动怎么办
策略:
新用户:热门/探索推荐
新物品:内容特征/规则召回
缓解:画像补充 + 探索机制
大数据支撑:
利用内容特征(不依赖行为)
冷启实验池8.2 实时性不足
表现:
用户刚看的物品没立即反馈
方案:
实时特征链路(Flink)
实时召回更新
特征新鲜度监控
权衡:
实时成本 vs 收益8.3 效果评估难
方法:
离线指标(AUC 等)
在线 AB 实验
业务指标(转化/时长)
注意:
离线好 ≠ 线上好
以 AB 为准九、小结
推荐系统 = 数据层(行为/特征)+ 召回层(多路候选)+ 排序层(粗排/精排)+ 重排层(业务约束),底层由大数据平台支撑:离线算模型与全量候选,实时算新鲜特征与召回。掌握链路分层与离线/实时分工,就能看懂一套完整的推荐系统是如何跑起来的。