Gorse Learning
本文最后更新于 2026年8月23日 下午
参考的资料:
入门
什么是推荐系统
三个要素:
- 记录行为
- 理解兴趣
- 预测可能喜欢
启动与使用
直接使用 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)

- 数据源输入输入层之后,传入检索层
- 检索层由多个推荐器构成,为用户生成候选物品
- 排序层合并来自不同推荐器的所有输出,删除用户已经看过的物品(已读物品),并根据用户与剩余物品互动的可能性对其进行评分
默认的管道是只推荐最新物品。
管道中的缓存
以下中间结果被缓存并定期更新:
- 用户到用户推荐器的用户邻居。
- 简单来说,用户 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 三层架构
- 各层独立扩展
- 无状态设计
✅ 负载均衡
- 一致性哈希
- 虚拟节点
- 最小化数据迁移
✅ 高可用
- 服务降级
- 熔断器
- 故障自动恢复
✅ 数据一致性
- 最终一致性
- 快速路径 + 慢路径
- 缓存保护