榜单模型与分布式任务调度

2023-01-27T14:11:02+08:00 | 17分钟阅读 | 更新于 2026-01-27T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 说清热榜业务的非功能性需求与难点,讲透 Hacknews / Reddit / 微博等典型热度算法的设计思路。
  2. 设计一套「定时批量计算 + 小顶堆维护 TopN + Redis ZSet 缓存」的高性能榜单方案。
  3. 落地「本地缓存 + Redis 缓存 + 数据库」的多级缓存,并理解「本地缓存用于容错兜底」这个高级技巧。
  4. 在 Go 里用 cron 跑定时任务,包括 cron 表达式、@every 语法和优雅退出。
  5. 解释分布式环境下定时任务重复执行的问题,并分别用 Redis 分布式锁和 MySQL 抢占式调度给出可落地的方案。

前置知识(如果下面任意一点生疏,先回看对应章):

  • 第02章 Gin + GORM:知道 DAO / Repository 分层、事务怎么写。
  • 第03章 Redis:知道 ZSet、过期时间、缓存读写基本套路。
  • 第07章 Kafka:知道消息是怎么从生产者到消费者的(后面本地缓存容错、修复思路会用到异步思维)。
  • 基本的 MySQL 索引概念(会看 KEYUNIQUE KEY)。

本章你会动手做的事

  • container/heap 写一个限容小顶堆,验证「堆满后只有比堆顶大才替换」的逻辑。
  • 给热榜计算服务套一层 Redis 分布式锁,把多实例重复计算压成一个实例算。
  • 用 MySQL 抢占式调度表,自己推演一次「节点崩溃 → 别的节点靠续约超时重新抢占」的全过程。

一、榜单业务需求分析

1.1 需求

展示一个热点榜单,例如 Top 50。从非功能性角度:

  • 榜单是首页或高频页面的核心模块,性能和可用性要求极高
  • 任何用户打开 APP 都会调用榜单接口

类比:热榜就像一个商场门口的「今日人气排行榜」大屏。所有人一进门就盯着看,所以它绝对不能卡、不能黑屏。但排行榜的内容是「算出来的」——不是用户一刷新就现场统计全场客流,而是后台提前算好贴上去的。这条「用户只看结果、后台提前算」的直觉,直接决定了我们后面要用「定时计算 + 缓存」而不是「实时计算」。

四个关键问题:

  1. 什么样的内容才算「热点」?
  2. 如何计算热度?
  3. 高并发下怎么保证性能?
  4. 榜单服务挂了怎么降低对系统的影响?

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.5P 直接取点赞数。

直觉:分子是「票数 - 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 批量计算流程

文章数可能非常多,一次性查全部会爆内存。采用批量处理:

  1. 从数据库分批拉文章(batchSize 比如 100)
  2. 查每批对应的点赞数,算 score
  3. 用小顶堆维护当前 Top 100
  4. 全部数据处理完,堆里就是 Top 100
  5. 写入 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 缓存)]
    end

9.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 设计通用的定时任务调度

  1. 数据库建一张表存待执行任务
  2. 所有实例都尝试「抢占」任务
  3. 抢到的执行

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 续约失败的查询条件(关键!)

「抢占」有两种合法情况:

  1. 没人调度status = 0 AND next_time <= now
  2. 调度方崩溃了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 上跑的服务。


十二、工程实践要点

  1. 小顶堆维护 TopN 比大顶堆省内存:固定容量,新值比堆顶大才替换。
  2. 批量分页拉数据,限制范围:只算最近 N 天的文章,老文章不可能上榜,省计算量。
  3. Redis 缓存过期时间 > 计算间隔:留够重试时间,避免缓存空窗。
  4. 本地缓存兜底是高级技巧:Redis 崩了用本地缓存顶,本地缓存过期时间故意设长一点(甚至永不过期)。
  5. 分布式锁加在 Job 而非 Service:业务逻辑本身没有「全局唯一」概念,是调度才有。
  6. 长锁 + 自动续约:避免每轮都抢锁,抢到的节点一直持有到退出。
  7. MySQL 抢占式调度必须配续约机制:否则节点崩溃任务永远卡在「执行中」。
  8. 抢占查询要用乐观锁:避免多实例并发抢同一个任务。
  9. 续约失败查询条件要 OR 进去:防止僵尸任务永远不被重新调度(这是作业二的考点)。
  10. 限流保护调度器:用 semaphore 限制并发抢占数量,防止内存打爆。

十三、面试加分:自研分布式任务调度平台

基于 MySQL 的实现就是一个简易的分布式任务调度平台。可在此基础上扩展:

  • 部门管理 + 权限控制:多团队共用一个调度平台
  • HTTP / gRPC 任务支持:调度任务 = 调用一个 HTTP 接口或 gRPC 方法,业务方无需接入 SDK
  • 任务执行历史:记录每次执行情况,方便审计和回溯
  • 告警通知:任务连续失败自动告警
  • 可视化界面:图形化配置 cron、查看运行状态

这套东西出去面试效果非常好。


自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 为什么热榜不能用「用户每次刷新实时全表扫描计算」?定时计算方案把成本从哪转移到了哪?
  2. 维护 TopN 为什么用小顶堆而不是大顶堆?堆已满时,新元素和堆顶是什么关系才会被保留?
  3. Redis 缓存的过期时间为什么要比「计算间隔」长好几倍?如果设成比间隔还短会出什么问题?
  4. Redis 短锁(每轮都抢)和长锁(启动抢一次一直持有)本质区别是什么?为什么老师选长锁?
  5. MySQL 抢占式调度里,抢占查询为什么必须有 OR (status=1 AND now-utime>阈值) 这个条件?漏掉会怎样?

动手练习(建议真做一遍)

  1. 写个限容小顶堆:实现 TopNHeap,用 1000 个随机 score 喂进去,验证最终 PopAll 返回的正好是 score 最大的 N 个,且顺序降序。
  2. 套一层 Redis 分布式锁:把示例的 DistributedRankingJob 接到你的 cron 任务上,起两个进程,观察日志确认同一时刻只有一个进程在算榜单。
  3. 推演崩溃恢复:在 MySQL 抢占式调度里,手动把某条 taskutime 改成「很久以前」,然后启动第二个实例,验证它能靠 OR 条件把这条僵尸任务抢回来执行。

十四、本章小结

  • 热榜算法核心三要素:用户行为 + 时间衰减 + 权重因子;Hacknews 公式 Score = (P-1) / (T+2)^G 是最简单的入门模型。
  • 性能方案核心:定时批量计算 + 小顶堆 TopN + 多级缓存;查询接口走「本地缓存 → Redis → 数据库」,并利用本地缓存做容错兜底。
  • 定时任务用 robfig/cron/v3,注意 WithSecondsStop 后等 ctx.Done 优雅退出。
  • 多实例下要解决重复执行:Redis 分布式锁适合「同一时刻唯一」,MySQL 抢占式调度适合做成通用调度平台。
  • MySQL 抢占式调度的两个关键设计:乐观锁防并发抢占 + 续约机制处理节点崩溃
  • 极致性能思路跳出代码本身:静态页 + OSS + CDN,甚至 APP 端定时拉取。

下一章(第10章)我们进入单体拆分微服务——当业务变复杂、团队变大,怎么把一整块单体安全地拆成独立服务,并为后面的不停机迁移、服务注册发现打基础。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!