学习目标
学完本章,你应该能够:
- 说清热榜业务的非功能性需求与难点,讲透 Hacknews / Reddit / 微博等典型热度算法的设计思路。
- 设计一套「定时批量计算 + 小顶堆维护 TopN + Redis ZSet 缓存」的高性能榜单方案。
- 落地「本地缓存 + Redis 缓存 + 数据库」的多级缓存,并理解「本地缓存用于容错兜底」这个高级技巧。
- 在 Go 里用 cron 跑定时任务,包括 cron 表达式、
@every语法和优雅退出。 - 解释分布式环境下定时任务重复执行的问题,并分别用 Redis 分布式锁和 MySQL 抢占式调度给出可落地的方案。
前置知识(如果下面任意一点生疏,先回看对应章):
- 第02章 Gin + GORM:知道 DAO / Repository 分层、事务怎么写。
- 第03章 Redis:知道 ZSet、过期时间、缓存读写基本套路。
- 第07章 Kafka:知道消息是怎么从生产者到消费者的(后面本地缓存容错、修复思路会用到异步思维)。
- 基本的 MySQL 索引概念(会看
KEY、UNIQUE KEY)。
本章你会动手做的事:
- 用
container/heap写一个限容小顶堆,验证「堆满后只有比堆顶大才替换」的逻辑。 - 给热榜计算服务套一层 Redis 分布式锁,把多实例重复计算压成一个实例算。
- 用 MySQL 抢占式调度表,自己推演一次「节点崩溃 → 别的节点靠续约超时重新抢占」的全过程。
一、榜单业务需求分析
1.1 需求
展示一个热点榜单,例如 Top 50。从非功能性角度:
- 榜单是首页或高频页面的核心模块,性能和可用性要求极高
- 任何用户打开 APP 都会调用榜单接口
类比:热榜就像一个商场门口的「今日人气排行榜」大屏。所有人一进门就盯着看,所以它绝对不能卡、不能黑屏。但排行榜的内容是「算出来的」——不是用户一刷新就现场统计全场客流,而是后台提前算好贴上去的。这条「用户只看结果、后台提前算」的直觉,直接决定了我们后面要用「定时计算 + 缓存」而不是「实时计算」。
四个关键问题:
- 什么样的内容才算「热点」?
- 如何计算热度?
- 高并发下怎么保证性能?
- 榜单服务挂了怎么降低对系统的影响?
1.2 热度模型:什么算热点
不同公司算法不同,但有共同规律:
- 综合用户行为:观看、点赞、收藏、评论、转发等
- 时间衰减:发布时间越久热度越低,避免老内容霸榜
- 权重因子:运营有意识控制(如加权、屏蔽、置顶)
理论上算法由产品经理定,研发只需要实现。但面试常考,要会几个典型算法。
flowchart LR
A[用户行为
点赞/收藏/评论/转发] --> H[热度算法]
B[发布时间] --> H
C[运营权重
加权/置顶/屏蔽] --> H
H --> R[最终 Score 排序]
R --> T[Top 50 榜单]二、典型热度算法
2.1 Hacknews 模型
Score = (P - 1) / (T + 2)^G
P= 投票数(点赞数)T= 发表以来小时数G= 重力系数(控制衰减速度,常用 1.5)
核心思想:得票数为主,时间衰减为辅。本课程 webook 就用这个模型,G = 1.5,P 直接取点赞数。
直觉:分子是「票数 - 1」,分母是「时间越久越大」的衰减项。一篇刚发 1 小时、100 赞的文章,和一篇发了 100 小时、200 赞的文章,谁更热?Hacknews 公式会让新内容有天然优势——这正是「避免老内容霸榜」的工程体现。
2.2 Reddit 模型
更复杂,考虑赞成票 + 反对票:
Score = log10(max(|x|, 1)) + sign(x) * (ts - 1134028003) / 45000
x= 赞成票 - 反对票y= 投票方向(赞多还是反多)z= 否定程度ts= 发帖时间(基准 2005.12.08)
核心:赞成票、反对票、发帖时间三者综合。
2.3 微博模型
微博没公布明确公式,只给了宽泛描述:「综合阅读量、讨论量、转发量、参与人数、参与次数、互动量等指标,按热度排序」。国内平台基本都这样。
2.4 webook 模型(本课程实现)
直接采用 Hacknews,最简化:
func score(likeCnt int64, utime int64) float64 {
// utime 是文章发表时间(毫秒)
hours := float64(time.Since(time.UnixMilli(utime)).Hours())
return float64(likeCnt-1) / math.Pow(hours+2, 1.5)
}
三、实时计算 vs 定时计算
3.1 实时计算的难点
如果想实时计算 Top 50:
- 全表扫描所有文章的点赞数和发布时间
- 每篇文章都算 score
- 全局排序找出前 50
难点:全表扫描 + 全局排序。 用户每次刷新榜单都这么干,数据库必崩。
3.2 解决方案:异步定时计算
每隔一段时间算一次榜单,把结果缓存起来。查询接口只读缓存。
要点:
- 间隔越短实时性越好,但消耗越大
- 计算时间可以稍长但不能太长(不能几个小时)
- 缓存设计要保证查询性能
- 要保证可用性,缓存挂了也要能返回数据
flowchart TD
subgraph 实时计算[每次刷新都算:不可行]
Q1[用户刷榜单] --> S1[全表扫描文章]
S1 --> S2[逐篇算 score]
S2 --> S3[全局排序取前 50]
S3 --> R1[返回]
end
subgraph 定时计算[后台周期算:本项目采用]
T[定时器触发] --> C[批量算 TopN]
C --> K[(写 Redis 缓存)]
Q2[用户刷榜单] --> R2[只读缓存返回]
end一句话定性:实时计算把「计算成本」摊到了每一次用户请求上,请求越多越崩;定时计算把成本固定成「每隔 N 分钟算一次」,用户请求只付「读缓存」的便宜账。
四、定时计算热榜
4.1 用 time.Ticker 实现简单定时器
func (s *RankingService) Start(ctx context.Context) {
ticker := time.NewTicker(time.Minute) // 每分钟算一次
defer ticker.Stop() // 退出时停止 ticker,防止泄漏
for {
select {
case <-ctx.Done():
return // 主动退出
case <-ticker.C:
_ = s.Rank(ctx) // 触发计算
}
}
}
注意点:
- 用 ctx 控制退出
- 用 select 监听
- 别忘了
ticker.Stop()
⚠️ 新手必踩的坑:忘了
ticker.Stop()。ticker 内部持有一个 channel 和底层定时器资源,goroutine 退出但没 Stop,定时器会一直跑、channel 一直堆积,造成 goroutine 与资源泄漏。凡是time.NewTicker,务必defer ticker.Stop()。
4.2 用 cron 库(更常用)
go get github.com/robfig/cron/v3
import "github.com/robfig/cron/v3"
c := cron.New(cron.WithSeconds()) // 启用秒级字段,必须设置
// 添加任务:AddFunc / AddJob 都是线程安全的
id, _ := c.AddJob("*/5 * * * * *", myJob)
c.Start() // 开始调度
// 优雅退出
ctx := c.Stop() // 只是暂停调度,正在运行的任务不会被中断
<-ctx.Done() // 等待所有任务结束
4.3 cron 表达式速查
格式(带秒):秒 分 时 日 月 周 [年]
0 */5 * * * * 每 5 分钟
0 0 * * * * 每小时整点
0 0 2 * * * 每天凌晨 2 点
0 0 0 1 * * 每月 1 号
便捷语法(推荐):
@every 1s 每秒
@every 1m30s 每 1 分 30 秒
@hourly 每小时
@daily 每天
老师原话:cron 表达式不小心就会写错,从来都是网上复制粘贴。需要时查文档:https://help.aliyun.com/document_detail/133509.html
⚠️ 新手必踩的坑:
cron.WithSeconds()不写会怎样。robfig/cron/v3 默认是 5 字段(分 时 日 月 周),如果你吭哧写了 6 字段*/5 * * * * *却没WithSeconds(),解析会直接报错或语义错乱。凡是带秒的 cron 表达式,必须cron.New(cron.WithSeconds())。
五、热榜算法实现
5.1 批量计算流程
文章数可能非常多,一次性查全部会爆内存。采用批量处理:
- 从数据库分批拉文章(batchSize 比如 100)
- 查每批对应的点赞数,算 score
- 用小顶堆维护当前 Top 100
- 全部数据处理完,堆里就是 Top 100
- 写入 Redis 缓存
flowchart TD
A[分批拉文章 batch=100] --> B[批量查点赞数]
B --> C[逐篇算 score]
C --> D{小顶堆
是否满?}
D -- 未满 --> E[直接入堆]
D -- 已满且新值>堆顶 --> F[弹堆顶+入新值]
D -- 已满且新值<=堆顶 --> G[丢弃]
E --> H{还有下一批?}
F --> H
G --> H
H -- 是 --> A
H -- 否 --> I[堆内即 Top100
写入 Redis]5.2 为什么用小顶堆维护 TopN
小顶堆:父节点值 ≤ 子节点值,堆顶是最小值。
算法:
- 堆大小固定 N(如 100)
- 新来一个元素:
- 堆未满 → 直接入堆
- 堆已满且新值 > 堆顶 → 弹出堆顶,新值入堆
- 堆已满且新值 ≤ 堆顶 → 丢弃
最终堆里就是 Top N。为什么不用大顶堆? 因为大顶堆要全部塞进去再取前 N,内存占用大;小顶堆只保留 N 个,省内存。
直觉:想象你在海选现场,只留 100 把椅子。来一个人,椅子没满就坐下;椅子满了,就把「当前最弱的那位」(堆顶)赶走、让新人坐。全程只维护 100 把椅子,不管来了几万人。这就是小顶堆省内存的本质。
5.3 优先级队列实现
package ranking
import "container/heap"
// ArticleScore 文章 + 热度分数
type ArticleScore struct {
Art domain.Article
Score float64
}
// priorityQueue 小顶堆,按 Score 升序
type priorityQueue []ArticleScore
func (p priorityQueue) Len() int { return len(p) }
func (p priorityQueue) Less(i, j int) bool { return p[i].Score < p[j].Score }
func (p priorityQueue) Swap(i, j int) { p[i], p[j] = p[j], p[i] }
func (p *priorityQueue) Push(x any) { *p = append(*p, x.(ArticleScore)) }
func (p *priorityQueue) Pop() any {
old := *p
n := len(old)
x := old[n-1]
*p = old[:n-1]
return x
}
// TopNHeap 限容的小顶堆
type TopNHeap struct {
h *priorityQueue
cap int
}
func NewTopNHeap(cap int) *TopNHeap {
pq := make(priorityQueue, 0, cap)
return &TopNHeap{h: &pq, cap: cap}
}
// Push 入堆:堆满则替换最小值
func (t *TopNHeap) Push(x ArticleScore) {
if t.h.Len() < t.cap {
heap.Push(t.h, x)
return
}
// 堆已满,比堆顶小则丢弃
if x.Score <= (*t.h)[0].Score {
return
}
// 比堆顶大,弹出堆顶,压入新值
heap.Pop(t.h)
heap.Push(t.h, x)
}
// PopAll 弹出全部并按降序返回
func (t *TopNHeap) PopAll() []ArticleScore {
res := make([]ArticleScore, 0, t.h.Len())
for t.h.Len() > 0 {
res = append(res, heap.Pop(t.h).(ArticleScore))
}
// 弹出是升序,反转成降序
for i, j := 0, len(res)-1; i < j; i, j = i+1, j-1 {
res[i], res[j] = res[j], res[i]
}
return res
}
5.4 计算服务实现(用 TDD 驱动)
package ranking
type RankingService interface {
RankTopN(ctx context.Context) error
}
type BatchRankingService struct {
articleSvc ArticleService
intrSvc InteractiveService
batch int // 每批大小
n int // TopN
scoreFunc func(like int64, utime int64) float64
cache RankCache // Redis 缓存
}
func (s *BatchRankingService) RankTopN(ctx context.Context) error {
// 步骤 1:初始化限容小顶堆
heap := NewTopNHeap(s.n)
offset := 0
// 只算最近 7 天的文章,老文章不可能上榜
since := time.Now().AddDate(0, 0, -7).UnixMilli()
for {
// 步骤 2:分批拉文章
arts, err := s.articleSvc.Batch(ctx, offset, s.batch, since)
if err != nil { return err }
if len(arts) == 0 { break }
offset += len(arts)
// 步骤 3:批量查点赞数
likeMap, err := s.intrSvc.BatchGetLike(ctx, "article", toIds(arts))
if err != nil { return err }
// 步骤 4:计算 score,维护 TopN
for _, art := range arts {
score := s.scoreFunc(likeMap[art.Id], art.Utime)
heap.Push(ArticleScore{Art: art, Score: score})
}
}
// 步骤 5:写入 Redis 缓存
top := heap.PopAll()
return s.cache.Save(ctx, top)
}
5.5 缓存设计要点
// 两个关键点:
// 1. 热榜数据通常不存数据库(除非要做大数据分析)
// 2. Redis 过期时间要比计算间隔长,留够重试时间
// 例如每分钟计算,过期时间设 10 分钟
字段设计:确保从缓存取出的数据就是查询接口需要返回的数据,字段一个不少。如果只需要标题 + 点赞数,就只存这些,避免再查一次文章详情。
⚠️ 新手必踩的坑:Redis 过期时间 ≤ 计算间隔。假设每分钟算一次、过期只设 30 秒,结果某次计算因为数据库抖动失败了,缓存到点被清空,榜单接口直接读不到数据返回空——这就是「缓存空窗」。正确做法:过期时间 = 计算间隔 × 数倍(如间隔 1 分钟设 10 分钟),给失败留重试余地。
六、Job 模块:组装定时任务
Job 不是 DDD 概念,但从项目结构上看,它和 Web、gRPC 同级,都是「对外暴露方式」。
6.1 Job 接口
package job
type Job interface {
Run(ctx context.Context) error
}
定义自己的接口方便后续用装饰器加监控、加日志、加分布式锁。
6.2 用 Builder 模式包装 cron.Job
type Builder struct {
p *prometheus.SummaryVec // 监控
l logger.Logger
}
func (b *Builder) Build(name string, j job.Job) cron.Job {
return cron.FuncJob(func() {
start := time.Now()
ctx := context.Background()
err := j.Run(ctx)
// 监控执行耗时
b.p.WithLabelValues(name).Observe(time.Since(start).Seconds())
if err != nil {
b.l.Error("任务执行失败", logger.String("name", name), logger.Error(err))
}
})
}
告警建议: 偶发失败不用管,连续失败或多次失败要看是什么问题。
6.3 启动与退出
func main() {
// ... wire 注入
c := cron.New(cron.WithSeconds())
c.AddJob("@every 1m", rankingJob)
c.Start()
// 优雅退出
stopCh := make(chan os.Signal, 1)
signal.Notify(stopCh, syscall.SIGINT, syscall.SIGTERM)
<-stopCh
ctx := c.Stop()
<-ctx.Done()
}
七、查询接口:极致性能的多级缓存
榜单查询接口是高并发 + 高可用接口,类似微博首页、小红书首页。
7.1 多级缓存方案
方案一: Repository 上直接组合 Redis 实现 + 本地缓存实现(推荐,开发快)
方案二: Cache 抽象层用装饰器组合两个实现(适合 Repository 有复杂逻辑时)
7.2 本地缓存实现
用原子操作即可,因为本质上是一个整体值,不需要 key-value 结构:
package localcache
import "sync/atomic"
// Value 泛型原子值,安全地存取整体数据
type Value[T any] struct {
val atomic.Value // 存 *T,避免泛型零值问题
}
func (v *Value[T]) Load() (T, bool) {
val := v.val.Load()
if val == nil {
var zero T
return zero, false
}
return *val.(*T), true
}
func (v *Value[T]) Store(val T) {
v.val.Store(&val)
}
7.3 Repository 组装
type CachedRankingRepository struct {
redis RankCache
local *localcache.Value[[]domain.Article]
exp time.Duration // 本地缓存过期
}
func (r *CachedRankingRepository) Get(ctx context.Context) ([]domain.Article, error) {
// 1. 优先读本地缓存
arts, ok := r.local.Load()
if ok {
return arts, nil
}
// 2. 本地没有,读 Redis
arts, err := r.redis.Get(ctx)
if err != nil {
return nil, err
}
// 3. 回写本地缓存
r.local.Store(arts)
return arts, nil
}
7.4 容错设计:本地缓存兜底
正常情况:本地缓存 → Redis → 数据库。
Redis 崩溃时的容错:
func (r *CachedRankingRepository) Get(ctx context.Context) ([]domain.Article, error) {
// 优先本地缓存
arts, ok := r.local.Load()
if ok { return arts, nil }
// 本地没有,尝试 Redis
arts, err := r.redis.Get(ctx)
if err != nil {
// Redis 崩了:再次尝试本地缓存,不检查过期
if arts2, ok2 := r.local.LoadForce(); ok2 {
return arts2, nil
}
return nil, err
}
r.local.Store(arts)
return arts, nil
}
flowchart TD
Q[查询榜单] --> L{本地缓存有?}
L -- 有 --> R1[直接返回]
L -- 无 --> RD{Redis 正常?}
RD -- 正常 --> C[读 Redis+回写本地]
RD -- 崩溃 --> LF{本地有
过期数据?}
LF -- 有 --> R2[返回过期数据兜底]
LF -- 无 --> E[返回错误]容错哲学:本地缓存「永不过期」的兜底思路,本质是用「稍微旧一点的榜单」换「服务不挂」。用户看到上一分钟的榜单,远比看到 404 强。这是可用性优先的经典取舍。
7.5 failover:节点本地缓存为空 + Redis 挂
最坏情况:节点刚启动没本地缓存,Redis 也崩了。让前端重试,下次请求大概率打到别的节点,别的节点可能有本地缓存。
7.6 Redis 缓存永不过期(终极容错)
让 Redis 缓存永不过期,定时任务每次刷新覆盖。这样:
- 数据库挂了 → 定时任务失败 → 但 Redis 旧数据还在
- 用户看到的是上一次的榜单,比 404 强
⚠️ 新手必踩的坑:本地缓存「永不过期」会不会脏数据爆炸?会。但它只存一份「整体榜单切片」,体积可控(Top 50 而已),且定时任务每轮覆盖。和「逐 key 缓存」不同,整体值缓存没有 key 膨胀问题,所以敢让它永不过期。
八、面试加分点
8.1 本地缓存 + Redis + 数据库三级缓存
- 查找:本地 → Redis → 数据库
- 更新:数据库 → 本地 → Redis(本地几乎不可能失败)
8.2 高级亮点
- 本地缓存预加载:启动时加载 / 快过期时提前刷新
- 本地缓存容错:本地缓存过期时间设长一点(比如正常 3 分钟,本地 5 分钟)。超过 3 分钟尝试刷新,刷新失败继续用「过期」数据。极端场景下本地缓存永不过期 + 异步刷新。
8.3 极致优化思路
- 算好榜单后直接生成静态页传 OSS + CDN
- 数据组成 JS 文件传 OSS + CDN
- 直接放到 Nginx 上
- APP 端定时拉数据本地缓存
极高并发下,Redis 也不一定扛得住。
九、分布式任务调度:问题与方案
9.1 问题:多实例重复执行
部署多个实例时,每个实例都会跑定时任务,同一个榜单被多次计算。
理论上结果一样,但浪费资源,要解决。
flowchart LR
subgraph 无锁[多实例无协调:浪费]
I1[实例1 计算榜单] --> D[(Redis 缓存)]
I2[实例2 计算榜单] --> D
I3[实例3 计算榜单] --> D
end
subgraph 有锁[抢锁后只有一个算:本项目采用]
L[谁抢到锁谁算] --> D2[(Redis 缓存)]
end9.2 方案一:Redis 分布式锁
思路:每个节点先抢锁,抢到的才计算。
go get github.com/gotomicro/redis-lock@latest
type DistributedRankingJob struct {
job job.Job
client *redislock.Client
key string
timeout time.Duration
}
func (d *DistributedRankingJob) Run(ctx context.Context) error {
// 尝试抢锁,过期时间 = 任务超时
locker := d.client.Obtain(ctx, d.key, d.timeout, nil)
if err != nil { return nil } // 别人拿到了,我跳过
defer locker.Release(ctx)
return d.job.Run(ctx)
}
关键决策:锁加在 Service 还是 Job 上?
老师选 Job。理由:计算热榜本身不存在「全局唯一」的概念,是 Job 调度才有「同一时刻只跑一次」的说法。
9.3 Redis 分布式锁的缺陷与改进
缺陷: 只能保证「同一时刻一个 goroutine 在计算」,但计算完释放锁后,下一轮别的节点又会计算。
改进:扩大锁的范围。 启动时抢锁,一直持有不释放(直到进程退出)。
type LongLockRankingJob struct {
job job.Job
client *redislock.Client
key string
timeout time.Duration
lock *redislock.Lock // 持有的锁
}
func (d *LongLockRankingJob) Run(ctx context.Context) error {
if d.lock == nil {
// 没拿到锁,尝试拿
lock, err := d.client.Obtain(ctx, d.key, d.timeout, nil)
if err != nil { return nil } // 别人有锁,跳过
d.lock = lock
// 拿到锁后开启自动续约
go d.lock.Refresh(ctx, d.timeout/2, nil)
}
// 自己持有锁,正常执行
return d.job.Run(ctx)
}
自动续约:锁快过期时自动延长,避免任务还没跑完锁就丢了。
优雅退出: 进程退出时调用 Close 主动释放。不调用也行,没人续约后过期自动释放。
⚠️ 新手必踩的坑:短锁下「每轮都抢、都算」。只看 9.2 的最短锁版本,每个实例每轮都
Obtain,抢不到就跳过——看似没问题,但 N 个实例每一轮都在抢锁、都在空跑Run的判断逻辑,纯属浪费。长锁(9.3)让「唯一计算者」长期持有,其余实例彻底不跑,资源最省。
十、方案二:MySQL 抢占式调度
10.1 思路
不用 Redis,用 MySQL 设计通用的定时任务调度:
- 数据库建一张表存待执行任务
- 所有实例都尝试「抢占」任务
- 抢到的执行
10.2 表设计
CREATE TABLE `task` (
`id` bigint UNSIGNED NOT NULL AUTO_INCREMENT,
`name` varchar(128) NOT NULL COMMENT '任务名',
`cron` varchar(64) NOT NULL COMMENT 'cron 表达式',
`status` tinyint NOT NULL COMMENT '0 待执行 1 执行中',
`next_time` bigint NOT NULL COMMENT '下次执行时间',
`utime` bigint NOT NULL COMMENT '最后续约时间',
`version` int NOT NULL DEFAULT 0 COMMENT '乐观锁版本号',
PRIMARY KEY (`id`),
KEY `idx_next_time` (`next_time`) -- 必须索引,加快抢占查询
);
10.3 状态流转
只有三种状态:待执行 → 执行中 → 待执行(释放后回到待执行)。
stateDiagram-v2
[*] --> 待执行
待执行 --> 执行中: 实例抢占成功
执行中 --> 待执行: 执行完释放
执行中 --> 待执行: 崩溃超时
被别的实例抢回10.4 续约问题
「抢占到了,但执行中崩溃了怎么办?」
答案:续约机制。 抢到任务的实例周期性更新 utime 证明自己还活着。其他实例发现 now - utime > 阈值 就认为该实例挂了,可以重新抢占。
sequenceDiagram
participant A as 实例A
participant DB as MySQL task 表
participant B as 实例B
A->>DB: 抢占 success,status=1
loop 每几秒续约
A->>DB: 更新 utime=now
end
Note over A: 突然崩溃
B->>DB: 发现 now-utime>阈值
B->>DB: 抢占该任务(执行中→待执行→执行中)10.5 Preempt 接口(Service 层)
type Scheduler interface {
Preempt(ctx context.Context) (*Task, error) // 抢占
Release(ctx context.Context, t *Task) error // 释放
Reset(ctx context.Context, t *Task, next time.Time) error // 重置下次时间
}
抢占用乐观锁:
func (d *TaskDAO) Preempt(ctx context.Context, t int64) (*Task, error) {
var task Task
err := d.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
// 步骤 1:找到符合条件的任务:状态为待执行 + 到了执行时间
err := tx.Where("status = ? AND next_time <= ?", 0, time.Now().UnixMilli()).
First(&task).Error
if err != nil { return err }
// 步骤 2:乐观锁更新状态为执行中
res := tx.Model(&Task{}).
Where("id = ? AND version = ?", task.Id, task.Version).
Updates(map[string]any{
"status": 1,
"version": gorm.Expr("version + 1"),
"utime": time.Now().UnixMilli(),
})
if res.RowsAffected == 0 {
return errors.New("抢占失败")
}
return nil
})
return &task, err
}
10.6 续约失败的查询条件(关键!)
「抢占」有两种合法情况:
- 没人调度:
status = 0 AND next_time <= now - 调度方崩溃了:
status = 1 AND now - utime > 续约超时阈值
完整查询:
now := time.Now().UnixMilli()
timeout := int64(30 * 1000) // 30 秒没续约视为崩溃
db.Where(
"(status = ? AND next_time <= ?) OR (status = ? AND ? - utime > ?)",
0, now, 1, now, timeout,
).First(&task)
⚠️ 新手必踩的坑:抢占查询漏掉 OR 第二种情况。很多实现只写
status = 0 AND next_time <= now,结果崩溃实例的任务永远停在status = 1,没人敢抢——成了「僵尸任务」,调度彻底卡死。这正是作业二的核心考点:必须把「崩溃超时」的OR条件加进去。
10.7 调度器设计
type MySQLScheduler struct {
db *gorm.DB
job job.Job
sem *semaphore.Weighted // 限制并发抢占数量
timeout time.Duration
}
func (s *MySQLScheduler) Schedule(ctx context.Context) {
for {
// 步骤 1:拿令牌(限制并发抢占数)
_ = s.sem.Acquire(ctx, 1)
go func() {
defer s.sem.Release(1)
// 步骤 2:抢占
t, err := s.Preempt(ctx)
if err != nil { return }
// 步骤 3:启动续约
keepAliveCtx, cancel := context.WithCancel(ctx)
go s.KeepAlive(keepAliveCtx, t)
// 步骤 4:执行任务
_ = s.job.Run(ctx)
// 步骤 5:释放
cancel()
_ = s.Release(ctx, t)
// 步骤 6:重置下次执行时间
_ = s.Reset(ctx, t, nextTime(t.Cron))
}()
}
}
限流原因: 不限制的话,极端情况下一次抢占几十万任务,内存爆掉。用 semaphore.Weighted 控制并发抢占数。
⚠️ 新手必踩的坑:不限制并发抢占数。调度循环是死循环,一旦任务表里有几十万条到期任务,
Schedule会瞬间 spawn 几十万 goroutine 去抢占,内存直接打爆。务必用semaphore.Weighted把并发数钉死在上限。
十一、其他调度框架
11.1 gocron
提供管理界面,本质是学一套新 API,没有难度。
11.2 K8s 任务调度
直接用 K8s 的 CronJob 资源,写一份 YAML 配置即可。适合已经在 K8s 上跑的服务。
十二、工程实践要点
- 小顶堆维护 TopN 比大顶堆省内存:固定容量,新值比堆顶大才替换。
- 批量分页拉数据,限制范围:只算最近 N 天的文章,老文章不可能上榜,省计算量。
- Redis 缓存过期时间 > 计算间隔:留够重试时间,避免缓存空窗。
- 本地缓存兜底是高级技巧:Redis 崩了用本地缓存顶,本地缓存过期时间故意设长一点(甚至永不过期)。
- 分布式锁加在 Job 而非 Service:业务逻辑本身没有「全局唯一」概念,是调度才有。
- 长锁 + 自动续约:避免每轮都抢锁,抢到的节点一直持有到退出。
- MySQL 抢占式调度必须配续约机制:否则节点崩溃任务永远卡在「执行中」。
- 抢占查询要用乐观锁:避免多实例并发抢同一个任务。
- 续约失败查询条件要 OR 进去:防止僵尸任务永远不被重新调度(这是作业二的考点)。
- 限流保护调度器:用 semaphore 限制并发抢占数量,防止内存打爆。
十三、面试加分:自研分布式任务调度平台
基于 MySQL 的实现就是一个简易的分布式任务调度平台。可在此基础上扩展:
- 部门管理 + 权限控制:多团队共用一个调度平台
- HTTP / gRPC 任务支持:调度任务 = 调用一个 HTTP 接口或 gRPC 方法,业务方无需接入 SDK
- 任务执行历史:记录每次执行情况,方便审计和回溯
- 告警通知:任务连续失败自动告警
- 可视化界面:图形化配置 cron、查看运行状态
这套东西出去面试效果非常好。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 为什么热榜不能用「用户每次刷新实时全表扫描计算」?定时计算方案把成本从哪转移到了哪?
- 维护 TopN 为什么用小顶堆而不是大顶堆?堆已满时,新元素和堆顶是什么关系才会被保留?
- Redis 缓存的过期时间为什么要比「计算间隔」长好几倍?如果设成比间隔还短会出什么问题?
- Redis 短锁(每轮都抢)和长锁(启动抢一次一直持有)本质区别是什么?为什么老师选长锁?
- MySQL 抢占式调度里,抢占查询为什么必须有
OR (status=1 AND now-utime>阈值)这个条件?漏掉会怎样?
动手练习(建议真做一遍):
- 写个限容小顶堆:实现
TopNHeap,用 1000 个随机 score 喂进去,验证最终PopAll返回的正好是 score 最大的 N 个,且顺序降序。 - 套一层 Redis 分布式锁:把示例的
DistributedRankingJob接到你的 cron 任务上,起两个进程,观察日志确认同一时刻只有一个进程在算榜单。 - 推演崩溃恢复:在 MySQL 抢占式调度里,手动把某条
task的utime改成「很久以前」,然后启动第二个实例,验证它能靠OR条件把这条僵尸任务抢回来执行。
十四、本章小结
- 热榜算法核心三要素:用户行为 + 时间衰减 + 权重因子;Hacknews 公式
Score = (P-1) / (T+2)^G是最简单的入门模型。 - 性能方案核心:定时批量计算 + 小顶堆 TopN + 多级缓存;查询接口走「本地缓存 → Redis → 数据库」,并利用本地缓存做容错兜底。
- 定时任务用
robfig/cron/v3,注意WithSeconds、Stop后等ctx.Done优雅退出。 - 多实例下要解决重复执行:Redis 分布式锁适合「同一时刻唯一」,MySQL 抢占式调度适合做成通用调度平台。
- MySQL 抢占式调度的两个关键设计:乐观锁防并发抢占 + 续约机制处理节点崩溃。
- 极致性能思路跳出代码本身:静态页 + OSS + CDN,甚至 APP 端定时拉取。
下一章(第10章)我们进入单体拆分微服务——当业务变复杂、团队变大,怎么把一整块单体安全地拆成独立服务,并为后面的不停机迁移、服务注册发现打基础。