本节把 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 与结构化并发 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。