《Go 语言编程实战》8.3 定时任务与延迟队列

TaskHub 还有一批时间驱动的任务:定时跑的统计与清理、延迟到某刻执行的通知与超时处理。本节用 time.Ticker 实测抖动、用分布式锁解决多实例重复触发、用 Redis ZSET 加 Lua 原子出队实现延迟队列,并对比各自的漂移与补跑语义。

本节把 TaskHub 推进到「时间驱动」:8.1、8.2 处理的是「收到事件就做」,本节处理「到点才做」——定时统计、延迟通知、任务超时,并解决多实例部署下定时任务被重复触发的难题。
适用版本:Go 1.27(实测 go1.27.0),Redis 7(redis:7-alpine 容器)。

8.3 定时任务与延迟队列

TaskHub 里有一类工作既不是「收到事件就做」,也不是「立刻做」,而是到某个时间点才做:每小时统计任务完成率、每天清理过期数据、任务分配后 30 分钟未开始则提醒、创建后 24 小时未完成则自动关闭。它们分两种:定时(周期性)和延迟(一次性、到点触发)。本节把两种都做对。

8.3.1 定时与延迟是两回事

类型触发方式TaskHub 例子实现
定时任务周期性、固定间隔每小时统计、每天清理time.Ticker 或 cron
延迟任务一次性、指定未来时刻30 分钟未开始提醒ZSET / 定时轮询

区别在于有没有明确的「下一个触发时刻」。定时任务是「每隔 T 执行」,延迟任务是「在 T 时刻执行一次」。混为一谈会导致实现选错——用定时轮询实现延迟任务会引入精度损失,用延迟队列实现定时任务会不断重复入队。

8.3.2 进程内定时器:time.Ticker

最直接的周期任务是 time.Ticker:

t := time.NewTicker(time.Hour)
defer t.Stop()
for range t.C {
	runHourlyStats(context.Background())
}

Ticker 保证「按固定间隔投递」,但间隔指的是相邻两次触发的间隔,不是绝对时钟对齐。实测每 50ms 触发一次的抖动:

$ go run ./ch8/cron
tick 0 at + 51ms drift=0s
tick 1 at +101ms drift=-7.291µs
tick 2 at +151ms drift=85.625µs
tick 3 at +201ms drift=-86.792µs
tick 4 at +250ms drift=-305.833µs

单次抖动在微秒级;tick 4 的 -305µs 属于调度抖动,Ticker 会按原始节奏重新对齐,不会累积漂移。更重要的是:Ticker 不补偿错过的触发。如果 runHourlyStats 跑了 70 分钟,下一小时的 tick 会被丢弃(Ticker 的 channel 容量为 1),你不会连补两次。这是很多「统计任务漏跑」的根因。

一次性的延迟用 time.Timer:

timer := time.NewTimer(80 * time.Millisecond)
<-timer.C // 80ms 后触发
$ go run ./ch8/cron
timer 80ms 触发

Timer 适合「本进程内、可随进程消亡」的短延迟(如连接超时、重试退避)。不适合跨进程的业务延迟任务——进程一重启,Timer 就没了。

Ticker 表达的是「每隔 T」,不是「每天凌晨 2 点」。后者需要 cron 表达式。用 github.com/robfig/cron/v3 这类库能把「0 2 * * *」翻译成「每天 02:00 触发」,并自带时区处理:

c := cron.New(cron.WithLocation(time.UTC)) // 显式时区,避免跨时区歧义
c.AddFunc("0 2 * * *", runDailyReport)      // 每天 UTC 02:00
c.Start()

(该库的 API 已在本机实测:AddFunc 返回 EntryID 与 error,Start/Stop 正常工作;上面用的是默认的 5 字段解析器,秒级调度需加 cron.WithSeconds()。)

「每隔一小时」和「每小时整点」是两种不同语义:前者相对上次执行时间累加,会随执行耗时漂移;后者对齐绝对时钟,不会漂移但可能因执行超时而跳过。要报表对齐整点就用 cron,要节流就用 Ticker,别用错。

8.3.3 多实例下的定时任务会重复触发

TaskHub 有 5 个实例,每个都跑 time.NewTicker(time.Hour) 的话,每小时统计会被跑 5 次。这是定时任务最经典的坑。三种解法:

方案做法特点
分布式锁触发前抢锁,抢到才跑简单,锁要带 TTL 防死锁
选主只有 leader 实例跑定时任务需要选主机制(如 Redis/etcd)
外部调度由 Kubernetes CronJob 或独立调度器触发最干净,但需额外组件

TaskHub 用分布式锁起步(复用 7.2 节的 SET NX EX + Lua 解锁)。每个实例都到点了,但只有一个能抢到锁真正执行:

func (s *Scheduler) runGuarded(ctx context.Context, name string, fn func()) {
	key := "cron:lock:" + name
	ok, _ := s.rdb.SetNX(ctx, key, s.instanceID, 90*time.Second).Result()
	if !ok {
		return // 别的实例已经在跑
	}
	defer s.unlock(ctx, key, s.instanceID) // Lua 校验令牌再删
	fn()
}

锁的 TTL 要略大于任务最长执行时间:太短会在任务没跑完时被别的实例抢走导致并发;太长则实例崩溃后恢复慢。

8.3.4 延迟队列:ZSET 方案

延迟任务需要一个「到点才取出」的队列。Redis 的 ZSET 天生适合:member 是任务,score 是执行时刻的 Unix 毫秒,取出时用 ZRANGEBYSCORE -inf now 拿所有到期的:

schedule := func(name string, at time.Time) {
	rdb.ZAdd(ctx, "jobs:delayed", redis.Z{
		Score:  float64(at.UnixMilli()),
		Member: name,
	})
}
schedule("job-A", now.Add(120*time.Millisecond))
schedule("job-B", now.Add(300*time.Millisecond))
schedule("job-C", now.Add(60*time.Millisecond))

轮询端按当前时间取到期的任务并执行,实测三条任务按各自时刻依次触发:

$ go run ./ch8/delay
排期 job-A 于 +120ms
排期 job-B 于 +300ms
排期 job-C 于 +60ms
触发 job-C 于 +54ms
触发 job-A 于 +110ms
触发 job-B 于 +313ms

job-C 排期 +60ms 却在 +54ms 触发,是因为轮询间隔 10ms,实际触发时刻取决于轮询精度。轮询间隔越短精度越高、但 Redis 压力越大,通常取 100ms~1s,取决于业务对延迟的容忍度。

8.3.5 原子出队:Lua 脚本

多实例同时轮询时,「取出」和「删除」必须原子,否则两个实例可能拿到同一条任务。用 Lua 把 ZRANGEBYSCORE 和 ZREM 合成一步:

popDue := redis.NewScript(`
local jobs = redis.call("ZRANGEBYSCORE", KEYS[1], "-inf", ARGV[1], "LIMIT", 0, 1)
if #jobs > 0 then
  redis.call("ZREM", KEYS[1], jobs[1])
end
return jobs`)

先取到期的最早一条,取到就立刻删——两步在同一个 Lua 脚本里执行,Redis 保证原子性,不会有两个实例抢到同一条。这和第 7.2 节的「比较令牌再删锁」是同一类手法:凡是「读—改」的复合操作,都要用 Lua 或事务包成原子。

拿到任务后,如果执行失败,要重新 ZADD 回队列(score 设为下次重试时刻),否则任务就丢了。这与 8.2 节的重试逻辑可以复用同一套退避函数。

原子性到底有多重要,用实验验证:4 个 worker 并发抢同一个 ZSET 里的 1000 条到期任务,每条任务用 LoadOrStore 检测是否被重复处理:

$ go run ./ch8/delaycon
4 个 worker 并发消费 1000 条延迟任务
处理总数=1000 重复=0

1000 条任务被 4 个并发 worker 精确消费,零重复、零遗漏。如果把 ZRANGEBYSCORE 和 ZREM 拆成两次独立的 Redis 命令,两个 worker 就可能在「A 取出、A 未删、B 取出」的窗口里拿到同一条任务——这正是 Lua 脚本要解决的竞态。判断标准很简单:凡是「先读后改」的复合操作,都要用 Lua、事务或 WATCH/MULTI 包成原子。

8.3.6 漂移、补跑与幂等

延迟任务有三个时间语义要明确:

  • 漂移:实际触发时刻晚于计划时刻(受轮询间隔、调度延迟影响)。上面的实验里 job-B 计划 +300ms、实际 +313ms,漂移 13ms。
  • 补跑:如果任务到点了但消费者当时不可用(进程重启),恢复后要不要补执行过去错过的任务?业务上要显式决定——「每天凌晨 2 点生成报表」错过要补,「每 5 分钟发心跳」错过就跳过。
  • 幂等:延迟任务同样可能被重复执行(重投、补跑),要像 8.2 节一样用 event_id 去重。

一个容易忽略的点:用 time.Now() 做计划时刻要小心时钟回拨。跨实例的延迟任务应该用统一的时钟源(Redis 服务器时间或数据库时间),避免各实例本地时钟不一致导致「有人觉得到点了、有人觉得没到」。TaskHub 的延迟任务统一以 Redis 的 TIME 命令为基准,不信任应用进程的本地时钟。

8.3.7 延迟队列的另一种实现:数据库表作队列

如果 TaskHub 已经重度依赖 Postgres,用一张表当延迟队列也是成熟方案。FOR UPDATE SKIP LOCKED 让多个 worker 各锁一行、互不阻塞:

DELETE FROM jobs WHERE id = (
    SELECT id FROM jobs WHERE run_at <= now()
    ORDER BY run_at
    FOR UPDATE SKIP LOCKED
    LIMIT 1
) RETURNING id;

同样的 4 个 worker 并发消费 1000 条,实测:

$ go run ./ch8/pgqueue
Postgres SKIP LOCKED: 处理=1000 重复=0

两种方案各有取舍:

方案优点缺点适用
Redis ZSET快、内存操作、天然按时刻排序数据在内存(可持久化但非事务性)量大、对速度敏感
Postgres 表与业务数据同库、可事务、可 SQL 查询依赖数据库、轮询有 IO 成本已有 Postgres、量中等

TaskHub 选 ZSET,因为它已经在用 Redis(7.2 节),不引入新的存储依赖。但如果延迟任务需要和业务数据在同一个事务里(比如「创建订单的同时排一个超时取消任务」),表作队列反而更合适——同库事务天然原子。选型的核心是看延迟任务与业务数据的一致性要求。

8.3.8 选型表

需求方案精度跨进程补跑
进程内短延迟time.Timer高否无
周期任务(单实例)time.Ticker高否不补
周期任务(多实例)Ticker + 分布式锁高是手动
延迟任务(跨进程)Redis ZSET + Lua轮询级是手动
复杂 cron 表达式cron 库 + 分布式锁秒级是手动
集群原生调度Kubernetes CronJob分钟级是可配

TaskHub 的选择:进程内短延迟用 Timer,周期任务用 Ticker + 分布式锁,业务延迟任务用 ZSET 延迟队列。Kubernetes CronJob 留到第 15 章部署时再评估。

8.3.9 常见坑

  • 多实例各跑一份 Ticker:定时任务被重复执行 N 次,必须加分布式锁或选主。
  • 锁 TTL 短于任务耗时:任务没跑完锁就过期,被别的实例并发执行。
  • ZSET 取出与删除非原子:两个实例抢到同一条任务,用 Lua 合并。
  • 轮询间隔设太短:空转压垮 Redis;设太长则延迟精度差,按业务容忍度取。
  • 任务失败不重新入队:延迟任务直接丢失,失败要 ZADD 回去。
  • 信任本地时钟:多实例时钟不一致会导致触发混乱,用统一时钟源。
  • 错过的任务不补跑也不告警:任务静默丢失,应该记录「错过」并决定补跑策略。
  • cron 表达式用了本地时区:跨时区部署时「凌晨 2 点」含义不同,显式指定时区。

小结

  • 时间驱动的任务分两种:周期性的定时任务用 Ticker,一次性的延迟任务用 ZSET。
  • Ticker 有微秒级抖动但不会累积漂移,且不补偿错过的触发;Timer 只适合进程内短延迟。
  • 多实例下定时任务会被重复触发,用分布式锁或选主解决,锁 TTL 要略大于任务耗时。
  • 延迟队列用 ZSET(member 任务、score 执行时刻),取出与删除必须用 Lua 保证原子。
  • 延迟任务要明确漂移、补跑、幂等三个语义,并用统一时钟源而非本地时钟。

到这里 TaskHub 的异步链路完整了:可靠投递、幂等消费、重试死信、定时与延迟。下一章回到并发本身——用 errgroup 把并发写得更结构化,用限流熔断保护下游,用 goleak 把泄漏挡在测试里。

阅读导航:上一节:8.2 幂等、重试与死信 · 下一节:9.1 errgroup 与结构化并发 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. 《Go 语言编程实战》目录
  2. 《Go 语言编程实战》18.3 上线、观测与迭代
  3. 《Go 语言编程实战》18.2 故障演练