Gorse Learning Note
本文最后更新于 2026年8月27日 下午
参考的资料:
入门
什么是推荐系统
三个要素:
- 记录行为
- 理解兴趣
- 预测可能喜欢
启动与使用
可以直接使用 docker 启动(参阅官方文档)。
基于 Web 的,直接访问 localhost:8088 可以浏览。
使用 curl 进行插入数据、创建物品、插入反馈、获取推荐等。
核心工作原理的概述
四个核心概念:
- 用户
- 基础信息:ID、标签等(如年龄、性别)
- 行为历史:浏览、点击、购买...
- 物品
- 基础信息:ID、标签
- 统计数据:热度、评分
- 反馈
- 用户 + 物品 + 类型 + 时间
- 推荐
- 根据历史行为预测用户可能喜欢的物品
流程图:
- 用户产生行为
- 系统记录反馈
- 模型定期执行分析、训练
- 生成推荐
- 用户看到推荐
- 循环迭代,回到 1
Gorse 的推荐策略:多源融合。
- 协同过滤:找到相似用户,推荐他们喜欢的物品
- 物品相似:推荐和用户历史物品相似的其他物品
- 热门推荐:推荐最热门的物品
- 最新推荐:字面意思
所谓融合策略:推荐结果 = 30% 协同过滤 + 30% 物品相似 + 20% 热门推荐 + 20% 最新推荐
Gorse 的架构设计
三层架构:
flowchart TD
A["用户 / 应用<br/>Web · App · 小程序"]
B["Server 节点<br/>RESTful API + 实时推荐"]
C["Master 节点<br/>模型训练 + 任务调度 + Dashboard"]
D["Worker 节点(多个)<br/>离线计算 + 批量推荐"]
E["存储层<br/>MySQL + Redis<br/>用户数据 + 物品数据 + 推荐缓存"]
A -->|HTTP / HTTPS| B
B -->|gRPC| C
C -->|gRPC| D
D --> E
关于 gRPC
RPC = Remote Procedure Call,远程过程调用。核心思想就是「像调用本地函数一样调用另一台机器上的函数」。
gRPC = Google 开源的一套 RPC 框架。
Protobuf = gRPC 通常使用的数据描述和序列化格式。你会写一个 .proto 文件定义「有哪些函数、参数是什么、返回值是什么」。
比如 Gorse 架构里:
Server ──gRPC──> Master
Master ──gRPC──> Worker假设 Master 提供一个函数:
GetModel(name string) ModelServer 想调用它。但问题是:Master 和 Server 是两个独立进程,甚至可能运行在不同机器上,Server 显然不能直接:
model := master.GetModel("ranking")gRPC 做的事情,就是让这种 远程调用看起来很像普通函数调用:
model, err := client.GetModel(ctx, request)实际上背后发生的是:
Server
│
│ 调用 GetModel(...)
↓
gRPC Client
│
│ 序列化成 Protobuf
│ 通过 HTTP/2 发送
↓
网络
↓
Master 上的 gRPC Server
│
│ 反序列化
↓
真正执行 GetModel(...)Master 节点
可以理解为 Gorse 架构的大脑。
- 模型训练
- AutoML:Automated Machine Learning(自动机器学习)
- 任务调度,触发 Worker
- Dashboard:监控、数据管理
Worker 节点
理解为 Gorse 架构的手脚。
- 批量推荐:为每个用户生成推荐列表
- 相似度计算:计算物品之间的相似,计算用户之间的相似度
- 水平扩展:启动多个 Worker,负载均衡
- 并行计算,每个 Worker 处理部分用户
Server 节点
理解为「嘴巴」。
- 提供 RESTful API
- 实时推荐
- 在线更新
- 无状态
综合
flowchart LR
A["用户行为"] --> B["Server"]
B --> C["DataStore<br/>MySQL"]
C --> D["Master<br/>定期加载数据"]
D --> E["训练模型"]
E --> F["Worker<br/>计算推荐"]
F --> G["CacheStore<br/>Redis"]
G --> H["Server"]
H --> I["返回用户"]
%% 局部调整布局
subgraph Train[" "]
direction TB
D --> E
end
subgraph Recommend[" "]
direction TB
F --> G
end
理解 Gorse 的管道(pipeline)
type Pipeline struct {
Config *config.Config // 配置
CacheClient cache.Database // Cache Store 客户端
DataClient data.Database // Data Store 客户端
Tracer *monitor.Monitor // 进度监控
Jobs int // 并行任务数
// 推荐模型
MatrixFactorizationItems *logics.MatrixFactorizationItems // 物品向量索引
MatrixFactorizationUsers *logics.MatrixFactorizationUsers // 用户向量
ClickThroughRateModel ctr.FactorizationMachines // CTR 预测模型
dontskipColdStartUsers bool // 是否跳过冷启动用户
}
- 数据源输入输入层之后,传入检索层
- 检索层由多个推荐器构成,为用户生成候选物品
- 排序层合并来自不同推荐器的所有输出,删除用户已经看过的物品(已读物品),并根据用户与剩余物品互动的可能性对其进行评分
默认的管道是只推荐最新物品。
管道中的缓存
以下中间结果被缓存并定期更新:
- 用户到用户推荐器的用户邻居。
- 简单来说,用户 A 的相似用户是 B C D,那么这个结果会被缓存
- 物品到物品推荐器的物品邻居。
- 和上面同理,相似物品缓存
- 非个性化推荐器的结果。
- 和具体用户没有太大关系的内容,比如说最新榜、热门榜
- 每个用户的排序器输出
- 就是最后排序器输出的结果
Gorse 工作原理
管道以 分布式 方式执行。Gorse 中有三种类型的节点:主节点、工作节点和服务器节点。也就是上面说的 Master Worker Server,不再赘述。
深入 Gorse 推荐系统:数据结构与存储层设计剖析
深入 Gorse 推荐系统:数据结构与存储层设计剖析 - 技术漫游 - 博客园
我感觉也是这个作者的 AI 文章自己洗了一遍...总之大概看看吧。
字符串-> 索引的映射
显然我们会有大量的 JSON(也可以是 Go 里面的 map[string][]),里面都是字符串作为 key。
字符串作为 key,效率很低(主要是哈希比较带来的开销是 $O(len(string))$
所以采用 ID <-> 索引 的双向映射。
具体而言,使用 FreqDict:
type FreqDict struct {
idToIndex map[string]int32 // ID → 索引
indexToId []string // 索引 → ID
frequencies []int // 频率统计
}下面的频率统计可以顺便用于统计最活跃的用户。
Dataset 稀疏矩阵
表示用户 - 物品的交互,如果使用二维数组,如 matrix := make([][]bool, 1000000),很浪费。
因此,采用稀疏存储:
type Dataset struct {
UserIndex *FreqDict // 用户字典
ItemIndex *FreqDict // 物品字典
UserFeedback [][]int32 // 用户 → 物品列表
ItemFeedback [][]int32 // 物品 → 用户列表
}存储双向的两份,空间换时间。
存储层
Gorse 支持很多数据库,比如 MySQL, PostgreSQL, MongoDB, Redis
实现方式是创建 Database 接口,然后各个数据库实现这些接口。
根据 配置项 | Gorse,
[database]
| 键 | 类型 | 默认值 | 描述 |
|---|---|---|---|
data_store |
string | 用于数据存储的数据库。 | |
cache_store |
string | 用于缓存存储的数据库。 | |
table_prefix |
string | 数据库中表的命名约定。 | |
cache_table_prefix |
string | table_prefix |
缓存存储数据库中表的命名约定。 |
data_table_prefix |
string | table_prefix |
数据存储数据库中表的命名约定。 |
DataStore 和 CacheStore
DataStore:存「原始业务数据」
- User, Item, Feedback
CacheStore:存「Gorse 计算出来的中间结果和推荐结果」
- User 的推荐列表
- 相似物品列表
这样,前者用关系型数据库,后者用 NoSQL 数据库如 Redis。
别的性能优化
批量加载
不使用从数据库对象一个一个 get,而是使用批量查询的 API,内部实现是从内部读取,减少对数据库的 I/O。
(其实就是 Redis 的最核心的应用吧。)
游标分页
要知道什么是游标分页,首先要知道什么是 offset 分页。
举个例子,比如表里有 1 万条数据,每次取 100 条。
SELECT * FROM feedbacks LIMIT 10000, 100;
-- 先找到前 10000 条,把它们全部跳过,再取后面的 100 条 --问题:前 10000 条虽然最后不要,但数据库还是得处理它们。
因此,越往后查,性能越差。
使用游标分页:每次查完,记录 last_id,下次查询的时候,使用 WHERE id > last_id,使用索引来直接定位。
并发控制
使用读写锁而不是一把大锁。
我觉得这算常识...
协同过滤算法深入
这个暂时空了,因为我没有很好的关于机器学习的基础。
Gorse 的分布式架构
先阐述一些基本概念、基本需求:
为什么需要分布式:
- QPS 不够(100 vs 10000)
- 内存不够(32GB vs 需要 100GB)
- 单点故障(挂了全挂)
Scale Out 水平扩展:
┌─────────┐ ┌─────────┐ ┌─────────┐
│ Server1 │ + │ Server2 │ + │ Server3 │
│ 100 QPS │ │ 100 QPS │ │ 100 QPS │
└─────────┘ └─────────┘ └─────────┘Scale Up 垂直扩展
┌─────────┐ ┌───────────┐
│ 32GB │ → │ 128GB │
│ 8 核 │ │ 32 核 │
└─────────┘ └───────────┘分布式采用 Scale Out 水平扩展。
整体架构
flowchart TD
U["用户请求"]
LB["Load Balancer<br/>Nginx / HAProxy"]
U --> LB
subgraph Servers["Server 集群 · 提供 API"]
direction LR
S1["Server1<br/>:8087"]
S2["Server2<br/>:8087"]
S3["Server3<br/>:8087"]
end
LB --> S1
LB --> S2
LB --> S3
M["Master<br/>:8086<br/>训练模型 · 调度任务"]
S1 -->|gRPC| M
S2 -->|gRPC| M
S3 -->|gRPC| M
subgraph Workers["Worker 集群 · 计算推荐"]
direction LR
W1["Worker1<br/>:8089"]
W2["Worker2<br/>:8090"]
W3["Worker3<br/>:8091"]
end
M -->|gRPC| W1
M -->|gRPC| W2
M -->|gRPC| W3
subgraph Storage["存储层"]
direction LR
DB["MySQL<br/>:3306"]
R["Redis<br/>:6379"]
end
W1 --> DB
W1 --> R
W2 --> DB
W2 --> R
W3 --> DB
W3 --> R
一致性哈希
其实就是实现负载均衡的策略。
举个例子,假设分配 User 交付于哪个 Worker 处理的策略是:
func getWorker(userId string) int {
hash := crc32.ChecksumIEEE([]byte(userId))
return int(hash) % numWorkers
}但是,假设增加一个 Worker,那么显然,原先的大部分 User 对应的 Worker 都不同了。
因此,我们需要采用 一致性哈希。
大致的原理是 哈希环:
想象有一个巨大的数字圆环,这里暂且假设为 0~99。
我们对 Users 和 Workers 同时进行哈希计算,哈希函数只会得到环上的数。
假设:
- hash(worker-A) = 20
- hash(worker-B) = 50
- hash(worker-C) = 80
并且
- hash(Alice) = 30
采用规则:从用户所在的位置顺时针走,碰到的第一个 Worker,就是负责这个用户的 Worker。
比如说 Alice 就该从 30 顺时针走,那么第一个就是 Worker-B 对应的 50。
考虑到增加 Worker 时,其实只会影响插入哈希的那一段,别的是不受影响的,这就是 哈希一致性。
Worker 的协同机制
Worker 启动时,主动连接 Master,注册自己,同步配置,拉取模型,开始工作。
Master 的任务分配:
- Master 定时扫描数据库,获得所有用户
- 如果有新用户,通知所有 worker 关于新用户的信息
- Worker 使用哈希值过滤,获得需要处理的用户
- 然后,并行计算,得到推荐列表
故障处理
Master 定时进行心跳检测 healthCheck()
- 如果 Worker 挂了,移除之
- 然后,一致性哈希重新分配(不在检测函数中处理),这样,这个 Worker 的用户会被分配给其他 Workers
Server 的水平扩展
无状态设计
每个 Server 都连接 cache Database,一般是 Redis。
然后,Server 是无状态的,体现在这个处理请求的函数:
// 处理推荐请求
func (s *Server) GetRecommendations(userId string, n int) []string {
// 1. 从缓存读取(无状态)
recommendations := s.cacheClient.GetRecommendations(userId, n)
if len(recommendations) >= n {
return recommendations[:n]
}
// 2. 缓存未命中,实时计算
return s.computeRecommendations(userId, n)
}为什么 Server 要无状态呢?因为 Server 随时可能挂掉,如果是有状态性的,如果某个用户的 Session 路由到了某个 Server 然后 Server 挂掉了,那么这个用户的 Session 就丢失了。
但是目前是任何修改直接写 Redis,这样挂了也没有影响,增减 Server 也不影响服务。
负载均衡
采用 Nginx。
数据一致性保证
最终一致性 / 旁路缓存策略 / 缓存更新问题
其实就是 MySQL 慢,而 Redis 快。(不准确,见下文 LLM)有时候可能读数据库得到推荐了,但是这个时候数据库已经开始写新数据了,然后用户先得到基于旧数据的推荐,新数据才更新进去。
解决方案:快速路径 + 慢路径
// Server 端
func (s *Server) InsertFeedback(feedback Feedback) {
// 1. 写入 MySQL(慢路径)
s.dataClient.InsertFeedback(feedback)
// 2. 快速更新:从推荐中移除已交互的
s.cacheClient.RemoveFromRecommendation(
feedback.UserId,
feedback.ItemId,
)
// 3. 触发增量更新(可选)
s.masterClient.TriggerIncrementalUpdate(feedback.UserId)
}引用
你没看懂很正常,因为这段文章把“数据一致性”讲得有点绕,而且那个“快速更新”没有先解释 为什么要删推荐结果。
先抓住一个非常具体的场景。
假设 Redis 里已经给用户 A 算好了推荐:
用户 A 的推荐列表: [手机,耳机,键盘,鼠标]这份推荐结果是 Worker 之前算好 放进 Redis 的。
现在用户 A 点了「耳机」。
于是系统收到一条 Feedback:
User=A Item=耳机 Type=click问题来了:
虽然 MySQL 现在会记录:
A 已经点击过耳机但是 Redis 里的推荐列表还是旧的:
[手机,耳机,键盘,鼠标] ↑ 还在这里如果用户立刻刷新推荐页,就可能又看到「耳机」。
这就很奇怪:
我刚刚都点过它了,怎么还继续推荐给我?
所以文章里的:
s.cacheClient.RemoveFromRecommendation( feedback.UserId, feedback.ItemId, )意思特别简单:
用户刚刚和某个 Item 发生了交互,那就立刻把这个 Item 从他当前缓存的推荐列表里删掉。
也就是:
原来的 Redis: A → [手机,耳机,键盘,鼠标] ↑ 点击了 立刻变成: A → [手机,键盘,鼠标]这就是所谓的:
快速更新:从推荐中移除已交互的
它并没有“重新计算推荐”。
只是做了一个特别便宜的小修补:
这个东西用户已经看过/点过了 ↓ 先别再推荐它
为什么叫“快速路径”?
因为真正重新计算推荐可能很慢。
假设用户点击耳机之后,我们想得到 完全最新 的推荐:
用户点击耳机 ↓ 写入 Feedback ↓ 重新读取用户行为 ↓ 重新跑推荐算法 ↓ 重新排序 ↓ 生成新的 Top N ↓ 写 Redis这个过程可能比较重。
你不可能要求用户:
“你稍等,我先重新跑一遍推荐算法,然后再给你刷新页面。”
所以采取一个折中方案:
用户点击耳机 ↓ 写 Feedback 到数据库 ↓ ┌────────┴────────┐ ↓ ↓ 快速路径 慢路径 立刻从 Redis 删除耳机 后台重新计算推荐 ↓ ↓ 几乎马上完成 可能过一会完成这就是文章所谓的:
快速路径 + 慢路径
具体一点。
快速路径:
RemoveFromRecommendation(user, item)只做:
[手机,耳机,键盘] ↓ 删除 [手机,键盘]非常快。
慢路径:
根据最新 Feedback 重新计算: A 现在喜欢什么? ↓ 重新生成推荐列表 ↓ [手机壳,蓝牙音箱,键盘,...]最终新推荐结果再覆盖 Redis。
所以你可以把这整个机制理解成:
先“止血”,再“彻底治疗”。
用户刚点击耳机:
立即做
不要再给他展示耳机。后台慢慢做
既然他刚刚点击了耳机, 那我是不是应该多给他推荐一些 蓝牙耳机、音箱、手机配件? 重新计算一遍。
另外我觉得你贴的文章里这句话:
“其实就是 MySQL 慢,而 Redis 快。”
不是特别准确。
真正的问题并不主要是:
MySQL 慢 Redis 快而是:
原始数据已经发生变化,但缓存里的推荐结果还是根据旧数据计算出来的。
这是一个典型的:
数据库里的事实 Redis 里的计算结果 A 点击了耳机 ←→ 还推荐耳机 新 旧两边暂时不一致。
甚至就算 MySQL 特别快,这个问题依然存在。
因为推荐结果本身通常是:
旧 Feedback ↓ Worker 计算 ↓ Redis 推荐列表新的 Feedback 出现以后,Worker 不一定瞬间重新算。
所以真正的问题是:
推荐缓存具有滞后性。
再来看你贴的三步代码就非常容易了:
func (s *Server) InsertFeedback(feedback Feedback) {用户产生了一条行为,比如:
A 点击 耳机第一步:
s.dataClient.InsertFeedback(feedback)永久保存:
MySQL: A 点击过 耳机这是事实数据。
第二步:
s.cacheClient.RemoveFromRecommendation( feedback.UserId, feedback.ItemId, )立刻修一下旧缓存:
Redis: A → [手机,耳机,键盘] ↓ 删除 A → [手机,键盘]第三步:
s.masterClient.TriggerIncrementalUpdate(feedback.UserId)意思可以粗略理解成:
“A 的行为变了,有空重新给 A 算一下推荐。”
之后可能生成:
A → [蓝牙音箱,手机壳,键盘,充电器]于是:
立即: 删掉已经交互的 Item 稍后: 生成真正更新后的推荐列表这就是整个思路。
你目前为了面试,其实记住下面这句话就很够用了:
Gorse 的推荐结果是预先计算并缓存的,所以用户产生新 Feedback 后,缓存里的推荐结果可能暂时过期。可以先快速从缓存中移除用户已经交互过的 Item,避免重复推荐;之后再由后台任务根据最新行为重新计算推荐结果。
这个思路比“因为 MySQL 慢 Redis 快”准确得多。
缓存穿透保护
缓存穿透是指客户端请求的数据在缓存中和数据库中都不存在,这样缓存永远都不会生效(只有数据库查到了,才会让 redis 缓存,但现在的问题是查不到),会频繁的去访问数据库。
Gorse 的解决方案:布隆过滤器。
布隆过滤器其实采用的是哈希思想来解决这个问题,通过一个庞大的二进制数组,根据哈希思想去判断当前这个要查询的数据是否存在,如果布隆过滤器判断存在,则放行,这个请求会去访问 redis,哪怕此时 redis 中的数据过期了,但是数据库里一定会存在这个数据,从数据库中查询到数据之后,再将其放到 redis 中。如果布隆过滤器判断这个数据不存在,则直接返回。这种思想的优点在于节约内存空间,但存在误判,误判的原因在于:布隆过滤器使用的是哈希思想,只要是哈希思想,都可能存在哈希冲突
代码
// 问题:恶意请求大量不存在的用户
func (s *Server) GetRecommendations(userId string) []string {
// 缓存未命中
recs := s.cache.Get(userId)
if recs == nil {
// ❌ 每次都计算,压垮系统
recs = s.compute(userId)
}
return recs
}
// 解决:布隆过滤器
type Server struct {
bloomFilter *BloomFilter
}
func (s *Server) GetRecommendations(userId string) []string {
// 1. 快速检查用户是否存在
if !s.bloomFilter.Contains(userId) {
return []string{} // 用户不存在,直接返回
}
// 2. 查询缓存
recs := s.cache.Get(userId)
if recs == nil {
recs = s.compute(userId)
}
return recs
}缓存雪崩
缓存雪崩是指在同一时间段,大量缓存的 key 同时失效,或者 Redis 服务宕机,导致大量请求到达数据库,带来巨大压力。
Gorse 的解决方案:设置随机的过期时间:
// 问题:大量缓存同时过期
func (s *Server) SetRecommendations(userId string, recs []string) {
// ❌ 所有缓存都是 1 小时过期
s.cache.Set(userId, recs, 1*time.Hour)
}
// 解决:随机过期时间
func (s *Server) SetRecommendations(userId string, recs []string) {
// 1 小时 ± 5 分钟
ttl := time.Hour + time.Duration(rand.Intn(600))*time.Second
s.cache.Set(userId, recs, ttl)
}缓存雪崩的解决方案,摘自 Redis 实战篇 | Kyle's Blog
- 给不同的 Key 的 TTL 添加随机值,让其在不同时间段分批失效
- 利用 Redis 集群提高服务的可用性(使用一个或者多个哨兵 (
Sentinel) 实例组成的系统,对 redis 节点进行监控,在主节点出现故障的情况下,能将从节点中的一个升级为主节点,进行故障转义,保证系统的可用性。 )- 给缓存业务添加降级限流策略
- 给业务添加多级缓存(浏览器访问静态资源时,优先读取浏览器本地缓存;访问非静态资源(ajax 查询数据)时,访问服务端;请求到达 Nginx 后,优先读取 Nginx 本地缓存;如果 Nginx 本地缓存未命中,则去直接查询 Redis(不经过 Tomcat);如果 Redis 查询未命中,则查询 Tomcat;请求进入 Tomcat 后,优先查询 JVM 进程缓存;如果 JVM 进程缓存未命中,则查询数据库)
容错与高可用
- 容错:系统某一部分出错时,整体还能继续工作,而不是一个地方出问题就全部崩掉。
- 高可用:系统尽量一直能对外提供服务,即使内部某些组件暂时故障。
- 服务降级:正常功能做不了时,退而求其次,提供一个“没那么好但还能用”的结果。比如个性化推荐失败了,就返回热门商品。
- Fallback(兜底):降级时实际返回的备用结果。比如热门推荐、默认推荐。
- 熔断器:当某个下游服务连续失败很多次后,系统暂时不再继续请求它,避免大量请求一直失败、拖垮整个系统。
- Closed(关闭):熔断器正常状态,请求正常通过。
- Open(打开):发现下游已经频繁出错,暂时禁止继续请求它。
- Half-Open(半开):过一段时间后,放几个请求进去试试看;成功了就恢复,失败了继续熔断。
服务降级
举个例子,Server 获取推荐的大致的代码逻辑是:
- 尝试从缓存获取
- 失败,则实时计算
- 实时计算失败,fallback 到获取预先缓存的热门推荐
代码
type Server struct {
circuitBreaker *CircuitBreaker
}
func (s *Server) GetRecommendations(userId string) []string {
// 1. 尝试从缓存获取
recs, err := s.cache.Get(userId)
if err == nil {
return recs
}
// 2. 缓存失败,尝试实时计算
if s.circuitBreaker.Allow() {
recs, err = s.computeRealtime(userId)
if err == nil {
return recs
}
s.circuitBreaker.RecordFailure()
}
// 3. 实时计算失败,降级到热门推荐
return s.getFallbackRecommendations()
}
// 降级策略
func (s *Server) getFallbackRecommendations() []string {
// 返回热门物品(预先缓存)
return s.cache.Get("popular_items")
}熔断器
感觉就是一个开闸关闸,关闸了之后还带检查,检查不过继续关,存在半开的中间状态。
代码
type CircuitBreaker struct {
state State // Open/Closed/HalfOpen
failureCount int
successCount int
failureThreshold int
timeout time.Duration
lastFailTime time.Time
}
func (cb *CircuitBreaker) Allow() bool {
switch cb.state {
case StateClosed:
return true // 正常状态,允许请求
case StateOpen:
// 熔断状态,检查是否到恢复时间
if time.Since(cb.lastFailTime) > cb.timeout {
cb.state = StateHalfOpen
return true // 尝试恢复
}
return false // 拒绝请求
case StateHalfOpen:
return true // 半开状态,允许部分请求
}
}
func (cb *CircuitBreaker) RecordFailure() {
cb.failureCount++
cb.lastFailTime = time.Now()
if cb.failureCount >= cb.failureThreshold {
cb.state = StateOpen // 打开熔断器
}
}
func (cb *CircuitBreaker) RecordSuccess() {
if cb.state == StateHalfOpen {
cb.successCount++
if cb.successCount >= 3 {
cb.state = StateClosed // 关闭熔断器,恢复正常
cb.failureCount = 0
}
}
}搭建分布式集群
主要是在 docker-compose.yml 里面配置。
当然,需要灵活扩展的话得上 k8s。
这个部分日后再单独讨论吧。
总结
✅ 分布式设计
- Master/Worker/Server 三层架构
- 各层独立扩展
- 无状态设计
✅ 负载均衡
- 一致性哈希
- 虚拟节点
- 最小化数据迁移
✅ 高可用
- 服务降级
- 熔断器
- 故障自动恢复
✅ 数据一致性
- 最终一致性
- 快速路径 + 慢路径
- 缓存保护
Gorse Server 工作原理详解
Server 在架构中的位置
- 用户请求传入 Server(无状态的服务)
- Server 的结果缓存到 Redis;通过 gRPC 同步到 Master
Server 的职责
- 对外提供 RESTful API
- 从 Redis 读取预先计算好的缓存结果
- 定期从 Master 同步配置和数据库连接信息
code
- 在
main.go创建 Server 实例 - 在
server.go实行初始化- 无状态
- 初始时不连接数据库,等待从 Master 获取配置
- 只需要知道 Master 的地址(其实是地址 + 端口)
- 在
server.go启动服务,通过 gRPC 连接到 master,使用go s.Sync()启动一个协程来同步,最后启动 HTTP 服务器Sync()跑一个无限循环,从 Master 请求元数据、解析- 若配置变了,就链接新的 Data Store / Cache Store
- 休眠一段时间,结束后继续循环
- 在
rest.go里面处理 RESTful API,注册完成后,具体的逻辑实现是一堆 CRUD。
走一遍完成的流程:
- 用户发起请求,比如是
GET /api/recommend/user123?n=10 - Server 的 HTTP Handler 收到请求,处理:
- 解析请求,得到参数
- 从 Redis 读取推荐的缓存
- 如果缓存命中,过滤已消费的物品
- 如果数量不够,采用降级策略(随机推荐等)
- 返回 JSON 响应
可以看出三个关键特性:
- 无状态:主要是负载均衡方便
- 配置热更新
- 降级策略
Server 性能优化要点
连接池
// Redis 连接池
redisClient := redis.NewClient(&redis.Options{
PoolSize: 100, // 连接池大小
MinIdleConns: 20, // 最小空闲连接
PoolTimeout: 4 * time.Second,
})
// MySQL 连接池
db.SetMaxOpenConns(100) // 最大连接数
db.SetMaxIdleConns(20) // 最大空闲连接
db.SetConnMaxLifetime(time.Hour)这里可以看出来连接池的两个要素:连接池大小 和 最小空闲连接。
缓存策略
过期时间:72h
更新频率:Worker 重新计算之后就会更新
对于热点数据:
- 流行物品,永久缓存
- 新物品:24h 缓存
采用批量操作
前面也说过,就是批量查询、批量加载。
Gorse Worker 架构详解
Gorse 的推荐计算引擎。
职责:
- 从 Master 同步配置和模型
- 生成离线推荐
- 写入缓存
Worker 结构体
代码
type Worker struct {
Pipeline
testMode bool
clickThroughRateModelId int64
// worker config
workerName string
httpHost string
httpPort int
masterHost string
masterPort int
tlsConfig *util.TLSConfig
cacheFile string
// database connection path
cachePath string
cachePrefix string
dataPath string
dataPrefix string
vectorPath string
vectorPrefix string
blobConfig string
blobStore blob.Store
vectorStore vectors.Database
// master connection
conn *grpc.ClientConn
masterClient protocol.MasterClient
latestCollaborativeFilteringModelId int64
latestClickThroughRateModelId int64
randGenerator *rand.Rand
// peers
peers []string
me string
// events
tickDuration time.Duration
ticker *time.Ticker
syncedChan chan struct{} // meta synced events
pulledChan chan struct{} // model pulled events
done chan struct{}
shutdown sync.Once
syncWait sync.WaitGroup
httpServer *http.Server
}这个地方东西有点多,挑几个我认为比较重要的字段:
-
Pipeline,和上文叙述过的 Pipeline 是一个东西,核心作用就是封装推荐计算的核心逻辑 -
自己的 HTTP 服务地址 + 端口
-
Master 地址与连接
-
Worker 节点列表、当前这个 Worker 在集群中的标识,用以一致性哈希
-
ticker 与通道,用于协程间的通信(Sync → Pull → Recommend)
A Ticker holds a channel that delivers “ticks” of a clock at intervals.
Worker 启动
首先主函数进行基本配置,然后启动 w.Serve()
而 w.Serve():
- 首先生成唯一的名称
- 创建进度跟踪器
- 用 gRPC 连接到 master
- 启动三大协程:Sync → Pull → Recommend
- Sync 用于同步配置
- Pull 用于拉取模型
- ServeHTTP 用于,字面意思
- 定义主循环
func loop():拉取用户、生成推荐、上报到 Master - 然后,开启一个无限循环,如果 channel 里面有定时触发 / 模型更新触发的信息,就调用
loop()
三大核心协程
Sync
大致执行:
一个无限循环:
- 从 Master 获取数据并更新
- 如果有新的就触发新的连接更新
- 检测模型版本更新
- 如果有新的就触发
Pull()协程
- 如果有新的就触发
- 更新集群信息用于一致性哈希
- 等待(睡眠一会儿)
Pull
收到 Sync 协程的信息时,循环继续(for range w.syncedChan)
- 拉取模型
- 如果有新模型,通过
w.pulledChan通知主循环
ServeHTTP
提供 Prometheus metrics 和健康检查 API
Prometheus 是一个开源的监控与告警系统。
Prometheus metrics 是程序专门暴露出来、给 Prometheus 采集的“运行状态数据”。
主要用于 k8s 里面检查。
用户推荐流程
详解 Pipeline.Recommend()
flowchart TD
A["开始推荐"] --> B["创建物品缓存"]
B --> C["并行处理每个用户<br/>p.Jobs 个并发"]
C --> D{"是否需要更新推荐?"}
D -->|"不活跃"| SKIP["跳过该用户"]
D -->|"缓存未过期"| SKIP
D -->|"冷启动用户"| SKIP
D -->|"需要更新"| E
E["① 协同过滤推荐<br/>获取用户向量<br/>在物品向量索引中搜索 TopK<br/>写入 CollaborativeFiltering"]
E --> F
subgraph Candidate["② 生成候选集"]
direction LR
F["调用多个 Recommender"]
F --> F1["ItemBased"]
F --> F2["UserBased"]
F --> F3["Latest"]
F --> F4["Popular"]
F --> F5["Custom"]
F1 & F2 & F3 & F4 & F5 --> F6["合并 · 去重"]
end
F6 --> G["③ 过滤候选集<br/>过滤已删除物品<br/>过滤隐藏物品<br/>添加替换候选(用户历史)"]
G --> H
subgraph Ranking["④ Ranking"]
direction LR
H{"排序方式"}
H --> H1["FM 模型<br/>CTR 预测"]
H --> H2["LLM 排序<br/>ChatGPT"]
H --> H3["不排序<br/>使用召回分数"]
end
H1 & H2 & H3 --> I["⑤ 应用替换衰减<br/>降低历史物品分数<br/>重新排序"]
I --> J["⑥ 写入缓存<br/>cache.Recommend<br/>cache.RecommendUpdateTime<br/>cache.RecommendDigest"]
J --> K["汇总所有用户"]
SKIP --> K
K --> L["上报 Metrics"]
L --> M["结束"]
其实就是前面 Pipeline 章节部分展开了说。
缓存过期策略
- 缓存为空
- 配置变化(digest 不匹配)
- 缓存超时(CacheExpire)
- 用户有新行为(activeTime > recommendTime)
负载均衡
分配用户的一致性哈希是在 Worker 架构里面实现的。
具体而言,分配用户的时候,采用上文所述的一致性哈希分配法。
三个架构之间的交互
flowchart TD
M["Master<br/>训练模型<br/>保存模型到 Blob Store<br/>提供配置、模型版本、Worker 列表<br/>接收 Worker 进度上报"]
W["Worker<br/>同步配置和模型版本<br/>下载模型<br/>生成用户推荐<br/>写入 Cache Store"]
C["Cache Store<br/>cache.Recommend<br/>cache.CollaborativeFiltering<br/>cache.RecommendUpdateTime"]
S["Server<br/>提供推荐 API<br/>读取预计算推荐<br/>应用实时过滤和补充"]
U["Client / User"]
M -->|"gRPC · GetMeta()"| W
W -->|"gRPC · PushProgress()"| M
W -->|"写入推荐结果"| C
C -->|"读取推荐结果"| S
S -->|"HTTP API"| U
不同架构之间都是通过 gRPC 交互的。
Worker 和 Server 不直接通信,通过 Cache Store 间接交互。
面向面试的部分
根据 LLM 提供的简历写法,进行对抗式的深挖,尽量做到准备充分,不被问穿。
Redis / 缓存部分
其实 Gorse 并没有像黑马点评那样专门去处理所谓缓存击穿、缓存穿透、缓存雪崩问题。
意思就是,那些博客园的文章,以及我对其所做的笔记的部分,都是扯淡,,,
不过也没辙,所以我给 AI 这样的 prompt:
但是,我总归还是需要 Redis 方面的亮点的,所以接下来你这样做:
- 检索互联网上是否有相关的优秀博客 / 资料,具体剖析 Gorse 在 Redis / 缓存处理这方面的优秀设计亮点的总结,如果有的话,就拿来作为我们博客上面写的亮点
- 如果没找到,结合源码自己归纳一下。
然后细讲 + 归纳,得到这样的结果:
缓存架构: 针对推荐结果计算成本高、传统缓存失效可能将重建压力传导至在线请求的问题,将推荐结果由 Worker 离线预计算并写入 Redis,通过用户数据更新时间、配置摘要和缓存有效期进行逻辑失效与后台刷新;结合 RediSearch 支持推荐结果的过滤与排序查询,在结果新鲜度、计算成本与在线请求延迟之间进行权衡。
但其实我也没太懂,所以让 LLM 针对这个写法给我细讲并展开。
首先理解推荐数据和普通业务数据的区别
对于普通的业务数据,比如 User ID:
- API 查询
- 假如 Redis Miss 了,再查数据库就行
但是对于推荐数据,是需要经过复杂的 ML 计算的。对于这一点,Gorse 的处理思路是后台 Worker 一直推荐计算存入 Redis,然后由 Server 读取 Redis 返回用户请求。
什么是 Cache Aside?
Cache Aside(旁路缓存)就是最常见的一种“数据库 + Redis”使用方式。
核心逻辑很简单:
读请求:
Client
↓
先查 Redis
↓
命中 → 直接返回
没命中
↓
查 MySQL
↓
把结果写入 Redis
↓
返回写请求通常是:
先更新 MySQL
↓
再删除 Redis 缓存为什么叫 Aside?因为缓存不在数据库内部,也不是数据库自动维护的,而是业务代码自己在旁边维护 Redis。
比如黑马点评查商户:
GET /shop/123
↓
Redis 查 shop:123
↓
没有
↓
MySQL 查 shop(id=123)
↓
写入 Redis
↓
返回这就是经典 Cache Aside。
而 Gorse 的推荐缓存和它不太一样。你的笔记里其实已经总结过这个区别:普通缓存更接近“数据库数据的副本”,而 Gorse 的推荐列表是经过推荐计算后得到的派生数据。
所以 Gorse 不能简单:
Redis miss
→ SELECT MySQL
→ 回填 Redis而更像:
DataStore
↓
推荐计算
↓
Worker
↓
CacheStore / Redis这也是为什么 Gorse 面对缓存击穿、雪崩时,设计思路会和《黑马点评》这种典型 Cache Aside 系统差很多。
缓存击穿问题
一般而言,应用 Redis 的项目如黑马点评里面的缓存击穿问题,是一个热点 Key 过期之后,所有请求都 Redis Miss,去请求数据库,导致击穿。
但是对于 Gorse,假设推荐结果失效之后,并不会强迫 Worker 立刻重新计算,而是 Server 先返回 Fallback 结果,然后 Worker 后台重新计算。
Gorse 不是用传统的“互斥锁防击穿”,而是通过 离线预计算 + 后台刷新,从架构上减少击穿问题。
但是,Worker 如何知道这个缓存是否应该重复计算?并非是 Redis 的某个 Key 的 TTL 过期之后就触发的。
在 Gorse 源码里面,有一个 checkRecommendCacheOutOfDate() 函数,大致检查这样的逻辑:
- 推荐列表是否空
- 用户配置是否有发生变化
- 指所请求的用户最近是否有产生新的行为,例如新的
click / like / favorite
- 指所请求的用户最近是否有产生新的行为,例如新的
- 距离上一次推荐计算是否已经超过 CacheExpire?
- 推荐配置的 digest 有没有变化?
- 所谓 digest 可以理解为推荐系统的配置做一个哈希函数之后得到的结果
- 所谓推荐系统的配置,就像是
cache_size之类的
RediSearch 相关
所谓 RediSearch:给 Redis Hash / Json 建立索引,支持条件过滤、排序、全文搜索、向量相似等。
在 Gorse 中的应用:推荐结果是一个类似文档的结构,包含如:
collection
subset
id
score
is_hidden
categories
timestamp建立起 User Key -> 对之的 K-V 结构之后,创建 RediSearch 索引。
这样,查询的时候,可以筛选条件过滤如 userid=123 && category = tech+ 按照 score 倒序再取 Top 20,而原生 Redis 做不到。
在结果新鲜度、计算成本与在线请求延迟之间进行权衡
- 要求结果永远最新鲜:计算成本巨大
- 计算成本低:结果不新鲜
- 在线实时计算:请求延迟高
向量召回
简历原文:
向量召回:针对大规模 Item 无法在在线请求中逐一计算相关性的问题,将文本等非结构化特征编码为 Embedding,利用向量相似度检索快速召回语义相关候选 Item,再交由后续排序流程生成最终结果,使向量检索负责缩小候选空间而非直接承担完整推荐决策。
所谓 Embedding 是什么其实是知道的,就是将特征转换为一个高维向量。
所谓「召回」:从海量的 Items 中根据 Embedding 的相似度快速找到一些相似 Items,构成候选集,之后的计算再从相对小规模的候选集中进行。
快速找到相似 Items 的部分,就是使用向量数据库完成的。
简历那部分的人话:东西太多,不能每次把所有东西都拿出来排序;所以先把 Item 内容转成向量,用向量相似度快速筛出一小批可能相关的内容,再精细排序。
拷打部分
用了这个 prompt:
引用
你是一名负责 后端开发实习 / 校招面试 的技术面试官。
我会给你:
- 我的简历项目描述;
- 必要的项目背景;
- 我可能提前整理的一些项目知识。
你的任务不是给我出题,而是进行一场尽可能接近真实后端技术面试的 项目深挖 / 简历拷打。
一、最重要的原则
- 你看不到我的代码
你必须始终假设:
面试官只有我的简历,以及我在面试过程中亲口告诉你的信息。
因此禁止因为我给你的背景材料中包含代码,就直接针对代码细节提问。
禁止:
- 问某个函数、类、变量是怎么写的;
- 问某一行代码为什么这么写;
- 问某个源码文件的实现;
- 根据我提供给 AI 的代码内容,假装面试官已经知道这些实现。
但是允许追问正常面试中候选人应该能够解释的实现机制,例如:
- 数据怎么流转;
- Redis Key / 数据结构为什么这么设计;
- MQ 消息怎么生产和消费;
- 幂等如何保证;
- 缓存和数据库如何保持一致;
- 并发问题如何处理;
- 一个请求经过哪些模块;
- 某个方案为什么采用这种实现。
判断标准:
一个真实面试官仅凭简历和候选人的口述,是否有理由提出这个问题?
如果没有,就不要问。
二、一次只问一个问题
严格模拟真实面试。
每轮只提出 一个主要问题。
不要一次输出:
- 为什么用 Redis?
- TTL 怎么设置?
- 如何保证一致性?
- Redis 挂了怎么办?
正确方式是:
为什么这里选择使用 Redis?
等待我回答。
如果我回答:
因为数据库查询比较慢,所以把推荐结果缓存到了 Redis。
下一题优先追问:
那你这个缓存什么时候失效?TTL 是怎么考虑的?
然后继续根据我的回答追问。
三、优先追问,而不是切换知识点
这是整个面试最重要的规则。
每次我回答以后,你首先判断:
这句话里面有没有值得继续追问的东西?
如果有,原则上 优先追问,不要马上换题。
一个重要技术点可以连续追问 3~5 层,必要时更多。
例如:
我:
为了降低数据库压力,我们使用 Redis 缓存推荐结果。
你:
为什么这个场景适合做缓存?
我回答以后:
那缓存的失效策略怎么设计?
继续:
更新推荐结果的时候,数据库和缓存之间怎么处理?
继续:
如果数据库更新成功,但是缓存更新失败呢?
继续:
那这种情况下你为什么接受最终一致性,而不是强一致性?
这才是正常的项目深挖。
四、追问来源
优先从我上一句话里寻找可以继续深入的点。
重点关注:
技术名词
如果我主动说:
- Redis
- MySQL
- PostgreSQL
- Kafka / RocketMQ
- gRPC
- RPC
- 分布式锁
- 一致性哈希
- Bloom Filter
- goroutine
- channel
- Elasticsearch
- Docker
这些词都可能成为新的追问入口。
但不要脱离项目突然变成八股考试。
例如我说:
用 Redis 缓存用户推荐结果。
可以问:
为什么这里选择 Redis?
不应该突然问:
Redis 有哪五种基本数据结构?
除非这个问题确实与当前讨论相关。
我的因果关系
如果我说:
因为 A,所以使用 B。
优先攻击这个因果关系:
- 为什么 A 会导致 B?
- 不使用 B 会怎样?
- B 真的是解决这个问题最合适的方式吗?
- 有没有其他方案?
- 为什么没选其他方案?
我的技术决策
如果我说:
我们使用消息队列进行异步更新。
可以继续:
为什么这里需要异步?
然后:
如果不用 MQ,直接同步调用有什么问题?
然后:
消息发送成功但是消费失败怎么办?
然后:
重试会不会产生重复消费?
然后:
那你的消费者怎么保证幂等?
我的指标
如果我说:
性能提高了 50%。
必须考虑追问:
这个 50% 是怎么测出来的?
如果我说:
可以支持高并发。
可以追问:
你实际测试过的 QPS 大概是多少?
如果数据听起来可疑,可以继续问。
不要默认接受简历上的数字。
我的模糊表达
例如:
做了一些优化。
立即问:
具体优化了什么?
或者:
提高了系统稳定性。
追问:
原来具体会发生什么问题,你做完之后有什么变化?
五、允许从项目自然延伸到八股
可以考察计算机基础,但必须存在一条自然的路径:
项目 → 技术 → 原理
例如:
项目使用 Redis
→ 为什么用 Redis
→ Redis 为什么快
→ 内存访问 / IO 多路复用 / 数据结构
这是合理的。
而不是:
项目使用 Redis
→ TCP 四次挥手是什么?
这种跳跃不要出现。
同理:
MQ
→ 重复消费
→ 幂等
→ 数据库唯一索引 / Redis / 状态机
合理。
MySQL
→ 查询慢
→ 索引
→ B+ 树
合理。
六、回答必须适合真实面试
我每个问题的目标回答时间是:
30~60 秒。
因此评价答案时,不只是判断“知识是否正确”,还要判断:
- 是否在 30~60 秒内能够讲完;
- 是否抓住重点;
- 是否一上来就回答问题;
- 是否讲了太多无关背景;
- 是否堆砌技术名词;
- 是否像背博客;
- 是否能够让面试官自然听懂。
不要要求我给出教科书式完整答案。
一个知识点理论上可以讲五分钟,不代表面试时应该讲五分钟。
优先训练:
先用 30~60 秒回答核心问题,剩余细节留给面试官追问。
如果我的答案明显过长,请记录下来,但真实面试过程中不要立刻长篇教学。
可以简单提醒:
这个回答有点长,真实面试建议压到一分钟以内。
然后继续面试。
七、不要主动帮我把答案补完整
如果我的回答有漏洞,不要马上告诉我标准答案。
例如:
我:
MQ 可以起到削峰的作用。
不要马上回答:
对,除此之外还有解耦、异步……
而应该继续问:
为什么它能削峰?
或者:
那峰值流量最后还是要被消费者处理,MQ 到底解决了什么问题?
我要训练的是被面试官追问,而不是让 AI 给我上课。
八、主动寻找“可攻击点”
我的每个回答都可能包含:
- 模糊表述;
- 缺少数据;
- 因果关系不充分;
- 技术选型说不清;
- 只知道结论不知道原因;
- 项目实际没做过;
- 只看过教程;
- 简历包装过度。
发现这些情况以后优先继续追问。
例如:
我:
为了保证高可用,我们使用 Redis。
可以立即追:
Redis 本身挂了怎么办?你这里说的“高可用”具体指什么?
九、场景变化型追问
当一个技术点已经解释清楚,可以适当改变约束继续追。
真实后端面试常见方式包括:
流量变化
如果 QPS 增长 10 倍呢?
组件故障
Redis 挂掉以后系统会怎样?
网络异常
MQ Broker 短暂不可用怎么办?
数据异常
消息重复了怎么办?
延迟异常
接口原来 50ms,现在突然变成 2s,你怎么排查?
容量变化
如果现在数据量增长 100 倍呢?
反事实
如果让你重新设计一次,你还会这么做吗?
但不要为了制造困难强行添加与项目无关的超大规模场景。
十、项目 ownership
重点判断:
这个项目到底是不是我真正理解并参与过的。
因此应该经常问:
- 这个模块为什么这么设计?
- 这个问题最开始是什么现象?
- 你当时怎么定位的?
- 为什么最终采用这个方案?
- 中间试过其他方法吗?
- 最大的问题是什么?
- 如果重新做你会改什么?
对于个人项目,不要强行问:
你们团队怎么分工?
而应该问:
哪部分是你重点改造或深入研究的?
十一、不要过度纠结“线上生产经验”
如果我是实习 / 校招候选人,尤其项目是个人项目或开源项目:
不要强迫我虚构:
- 百万 QPS;
- 真实生产事故;
- 大规模线上用户;
- 公司级 SLA。
如果没有真实线上数据,可以问:
你有没有做过压测?
如果没上线,你怎么判断这个设计能够解决这个问题?
这个优化的收益有没有通过 benchmark 或实验验证?
重点判断工程思维,而不是逼候选人编造数据。
十二、问题难度
目标职位:
后端开发实习 / 校招。
问题应该接近真实一面 / 二面。
重点:
- 项目理解;
- 后端基础;
- 技术选型;
- 数据库;
- Redis;
- MQ;
- 并发;
- 网络;
- 分布式基础;
- 性能和可靠性;
- 场景设计。
不要默认按照高级后端工程师 / 架构师标准要求。
但是,如果我主动在简历中写了一个高级概念,就可以针对这个概念认真追问。
原则:
简历写得越深,面试官就越有理由问得深。
十三、维护“追问树”
面试过程中,你在内部维护一个追问树。
例如:
Redis 缓存推荐结果 ├─ 为什么缓存 ├─ Key / Value ├─ TTL ├─ 更新策略 │ ├─ Cache Aside │ ├─ 更新 DB 失败 │ └─ 删除缓存失败 ├─ 缓存击穿 ├─ Redis 故障 └─ 10 倍流量
只要当前分支还有明显价值,就优先沿当前分支继续。
只有出现以下情况才切换:
- 当前知识点已经确认掌握;
- 我明显不会;
- 已经连续追问足够深入;
- 问下去已经明显超出实习 / 校招要求;
- 当前问题开始脱离项目。
然后再返回之前尚未探索的重要分支。
十四、真实面试模式
正常情况下不要输出:
- 标准答案;
- 大段点评;
- 下一题为什么这么问;
- 当前评分;
- 完整追问树;
- 教学内容。
你的正常输出应该非常接近真实面试官:
为什么这里选择消息队列,而不是直接同步调用?
等我回答即可。
十五、复盘模式
只有当我输入:
复盘
才暂停面试。
此时输出:
- 本轮整体评价
按照真实后端实习 / 校招标准评价。
- 问题逐题复盘
对于重要问题给出:
- 我的回答哪里好;
- 哪里有问题;
- 面试官可能继续攻击哪里。
- 30~60 秒参考回答
对于我答得不好的核心问题,给一个:
适合真实面试口述的 30~60 秒答案。
不要写成长篇教程。
- 知识状态
分成:
- 已经比较稳;
- 能回答但容易被追问击穿;
- 明显不会;
- 需要补充项目事实。
- 简历风险
指出:
- 哪句话容易被拷打;
- 哪句话夸大;
- 哪句话缺乏指标;
- 哪句话我实际上解释不清楚。
- 下一轮重点
列出下一轮最值得继续练的 3~5 条追问链。
十六、开始方式
收到我的项目资料以后:
先理解项目,但不要把我提供给 AI 的全部材料视为“面试官已知信息”。
真实面试官初始只知道:
我简历上实际写出来的内容。
其他材料只用于你判断我的答案是否合理以及设计后续追问。
然后直接开始:
请你先简单介绍一下这个项目,重点讲一下它解决什么问题,以及你主要做了什么。
从此进入真实面试模式。
一次只问一个问题。
一些面试题
明白,那我换个问法:既然算法本身不是你设计的,你为什么还需要把整个推荐服务拆成 Master、Worker、Server 三类节点?这个架构主要是在解决什么工程问题?
拆开来一个是解决推荐计算和在线请求之间的延迟需求不一样的问题。
- 推荐计算要求重,而且可以离线执行
- 在线接口要求相应快和尽量稳定
因此,这样各个职责解耦,各司其职而且每个节点扩容也方便,也避免别的服务的瓶颈拖垮另外的服务。
如果同一个用户的推荐缓存已经失效,短时间内又同时进来很多个 Server 请求,这些请求都发现 Redis 里没有可用结果,然后都去通知 Master。你怎么避免同一个用户的推荐任务被重复调度、多个 Worker 重复计算?
这个解决方案其实比较简单,架构是 Worker 周期性生成所有用户的离线推荐,Server 只读缓存和 Fallback 结果。并没有 Server Cache miss 之后通知 Master 调度重算的逻辑。