Go 基础体系 · 第 76/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go 定时任务:robfig/cron、时区、防重入与分布式调度
本文以 Go 1.26.4 和稳定版 github.com/robfig/cron/v3 v3.0.1 为基准。robfig/cron 是进程内调度器:它解析计划、计算下一次时间并在 goroutine 中调用 Job。它不持久化执行记录,不在进程离线时自动补跑,也不自动解决多副本重复执行。理解这些边界,比记住 AddFunc 更重要。
1. 内部模型:Cron、Entry、Schedule 与 Job
Cron 拥有 entry 集合、控制 channel 和一个调度循环;Entry 保存 ID、Schedule、Job、上次与下次时间;Schedule.Next(time.Time) 计算严格晚于输入的下一次触发;Job.Run() 执行业务。函数通过 FuncJob 适配为 Job,wrapper 再装饰其并发或恢复语义。
parser -> Schedule.Next(now) -> Entry.Next
scheduler goroutine
| timer fires
+-> go Job.Run()
+-> recompute Entry.Next
调度循环不等待 Job 完成,因此默认允许同一 Entry 重叠。时钟到点只代表 goroutine 被安排启动,不代表准时完成。进程重启会重新计算未来的 Next,停机期间的触发不会形成持久 backlog。
2. 安装、版本与最小可运行程序
模块锁定版本,避免 v1/v2 示例与 v3 API 混用:
go mod init example.com/cron-demo
go get github.com/robfig/cron/v3@v3.0.1
go mod tidy
go test ./...
location, err := time.LoadLocation("Asia/Shanghai")
if err != nil {
return fmt.Errorf("load schedule location: %w", err)
}
scheduler := cron.New(cron.WithLocation(location))
id, err := scheduler.AddFunc("0 3 * * *", func() {
logger.Info("cleanup started")
})
if err != nil {
return fmt.Errorf("add cleanup schedule: %w", err)
}
logger.Info("schedule registered", "entry_id", id)
scheduler.Start()
defer scheduler.Stop()
生产任务不能忽略回调中的错误。AddFunc 没有 error 返回通道,应让闭包调用返回 error 的业务函数,并在任务运行边界记录一次结果。
3. 五字段、秒字段和 descriptor
默认 parser 使用分钟、小时、月日、月份、周日五字段,0 3 * * * 表示每天 03:00。秒字段不是默认行为;若配置采用六字段,必须显式启用 cron.WithSeconds(),否则启动即报错。不要根据字段数量自动猜格式,这会让一次配置变更改变全部计划含义。
secondsCron := cron.New(
cron.WithSeconds(),
cron.WithLocation(time.UTC),
)
if _, err := secondsCron.AddFunc("*/15 * * * * *", heartbeat); err != nil {
return fmt.Errorf("add heartbeat schedule: %w", err)
}
需要精确控制时直接创建 parser:
parser := cron.NewParser(
cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow | cron.Descriptor,
)
schedule, err := parser.Parse("@every 10m")
@every 10m 表示固定间隔,不等于每个墙上时钟的整十分钟。配置文档应公布接受的方言,并在加载阶段解析全部表达式,不能等到服务运行后才发现坏配置。
4. 日期字段与 Next 语义
数字、范围、列表、步长和星号组合成计划。月日与周日同时受限时,cron 传统语义容易被误解,必须用具体日期测试。Schedule.Next(t) 返回晚于 t 的下一次时间,适合预览、启动校验和补跑规划。
func NextRuns(spec string, loc *time.Location, from time.Time, n int) ([]time.Time, error) {
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
schedule, err := parser.Parse(spec)
if err != nil {
return nil, fmt.Errorf("parse cron schedule: %w", err)
}
runs := make([]time.Time, 0, n)
next := from.In(loc)
for range n {
next = schedule.Next(next)
runs = append(runs, next)
}
return runs, nil
}
管理界面应显示接下来若干次带时区的时间,而非只回显表达式。这能在发布前发现“周一”编号、月末和时区理解错误。
5. Location、CRON_TZ 与夏令时
WithLocation 给整个 Cron 默认时区,表达式也可通过 parser 支持的 CRON_TZ=Area/City 指定单条计划。使用 IANA 名称,例如 Asia/Shanghai、Europe/Berlin;CST 有歧义,固定 UTC+8 又无法表达夏令时规则。
容器的 tzdata 可能缺失。Go 可通过系统 zoneinfo 或嵌入 time/tzdata 提供数据库;生产镜像必须验证。夏令时切换会产生不存在的当地时间或重复时间:要求“每天当地 02:30”的业务必须决定跳过、顺延还是执行两次。库按时间规则计算,但产品语义不能交给库猜。
loc, err := time.LoadLocation("Europe/Berlin")
if err != nil {
return fmt.Errorf("load Europe/Berlin: %w", err)
}
runs, err := NextRuns("30 2 * * *", loc,
time.Date(2026, 3, 27, 0, 0, 0, 0, loc), 5)
对结算日等关键业务,使用业务日历生成待执行窗口通常比纯 cron 表达式更清楚。
6. AddFunc、AddJob、Entry 与 Remove
AddFunc 注册函数,AddJob 注册实现 Run 的对象;两者返回稳定的 EntryID。Entry(id) 返回快照,可查看 Prev、Next、Valid;Entries() 用于诊断;Remove(id) 阻止未来触发,但不会取消已经启动的 Job。
type CleanupJob struct{ store *Store }
func (j *CleanupJob) Run() {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
if err := j.store.Cleanup(ctx); err != nil {
j.store.Logger().Error("cleanup failed", "err", err)
}
}
id, err := scheduler.AddJob("0 3 * * *", &CleanupJob{store: store})
if err != nil {
return err
}
entry := scheduler.Entry(id)
if !entry.Valid() {
return errors.New("cleanup entry is invalid")
}
动态替换计划时先成功解析并注册新 entry,再移除旧 entry,并保存映射。这个过程不是跨进程事务;关键调度配置更适合外部控制面或持久化任务系统。
7. Start、Run、Stop 的生命周期
Start 在后台启动调度循环并立即返回;Run 占用当前 goroutine;重复 Start 不应成为业务控制手段。Stop 停止调度新 Job,并返回一个 context:其 Done 在当时运行中的 Job 结束后关闭。它不会把取消信号传给 Job,因为 Job.Run 没有 context 参数。
scheduler.Start()
<-processCtx.Done()
waitCtx := scheduler.Stop()
select {
case <-waitCtx.Done():
return nil
case <-time.After(30 * time.Second):
return errors.New("cron jobs did not stop within 30s")
}
要主动取消 Job,需由应用持有进程 context,闭包为每次执行派生 timeout。停止顺序是:取消根 context、Stop 阻止新触发、等待 Stop context,再关闭数据库等依赖。不能先关连接再等待任务。
8. 为每次执行建立 context、run ID 与预算
cron 的函数签名不带 context,可以用工厂闭包注入进程生命周期。每次触发创建独立 timeout 和 run ID,避免多个执行共享可变 context。
func timedJob(parent context.Context, timeout time.Duration, run func(context.Context) error) func() {
return func() {
ctx, cancel := context.WithTimeout(parent, timeout)
defer cancel()
started := time.Now()
err := run(ctx)
logger.InfoContext(ctx, "cron run finished",
"duration", time.Since(started), "err", err)
}
}
_, err := scheduler.AddFunc("*/5 * * * *",
timedJob(processCtx, 4*time.Minute, refreshIndex))
SQL 用 ExecContext,HTTP 请求用 NewRequestWithContext。timeout 后外部操作可能已提交,因此有副作用任务仍需幂等。日志里的 context 取消错误要区分进程关停、单次 deadline 和下游失败。
9. 默认重叠、Delay 与 Skip
默认每次触发都在新 goroutine 执行,前一次未结束时会重入。DelayIfStillRunning 串行等待,保留每次触发,但延迟可能无限累积;SkipIfStillRunning 在忙时跳过,适合只需最新一次刷新的任务。wrapper 顺序会改变观察与恢复边界,应集中配置。
scheduler := cron.New(
cron.WithLocation(location),
cron.WithChain(
cron.Recover(cron.DefaultLogger),
cron.SkipIfStillRunning(cron.DefaultLogger),
),
)
Delay 的等待 goroutine 不接收 context,长时间积压不是可靠队列。若每次计划都必须最终执行,应把 cron 仅作为触发器,把带时间窗口幂等键的任务写入 Asynq 或数据库队列表。Skip 也不是失败,必须有 skipped 指标,否则运营只能看到“没有报错但数据没刷新”。
10. Recover、错误处理与重试
v3 默认不自动恢复 panic;cron.Recover(logger) 可避免一个 Job panic 终止进程。panic 表示缺陷,恢复后仍需堆栈、告警和修复。普通 error 不应 panic,Job 的运行边界记录一次并更新执行表。
robfig/cron 不提供业务重试。可在一次 Job 内做有上限、受总 deadline 控制的瞬时重试,但更可靠的方式是触发持久队列。重试必须分类:网络暂时失败可退避,参数非法不可重试,未知提交状态先按幂等键查询。
func enqueueWindow(ctx context.Context, day time.Time) error {
key := "daily-report:" + day.Format("2006-01-02")
_, err := client.EnqueueContext(ctx, makeReportTask(day),
asynq.TaskID(key), asynq.MaxRetry(10))
if errors.Is(err, asynq.ErrTaskIDConflict) {
return nil
}
return err
}
这里 Cobra/Viper 不参与运行语义;计划配置可由 Viper 加载,但调度器仍必须自己管理验证与生命周期。
11. 业务幂等、防重入与执行账本
单进程 Skip 只能防当前 Cron 实例重入,进程重启或多副本仍会重复。按计划窗口生成稳定键,例如 job_name + scheduled_date,在数据库建立唯一约束。事务中插入 run 记录,冲突表示已领取;完成后写状态和结果摘要。
CREATE TABLE job_runs (
job_name text NOT NULL,
window_start timestamptz NOT NULL,
status text NOT NULL,
owner_token bigint NOT NULL,
started_at timestamptz NOT NULL,
finished_at timestamptz,
PRIMARY KEY (job_name, window_start)
);
唯一记录解决“同窗口只认一个结果”,不一定取消旧执行。业务更新还应带状态条件或 fencing token,拒绝过期 owner 的写入。不要只用内存 mutex,它对其他副本和重启无效。
12. 多副本、租约与 fencing
服务部署 N 个副本,每个进程内 cron 都会触发 N 次。可选策略是:明确允许每副本执行;只部署一个 scheduler 副本;使用数据库/协调系统选主;或把触发交给 Kubernetes CronJob、Asynq Scheduler、云调度器。
分布式锁必须有租约,避免持有者崩溃永久占锁;但租约过期不代表旧进程已经停止。网络暂停后旧进程可能恢复并继续写,因此锁服务返回单调递增 fencing token,下游只接受不小于当前 token 的写入。没有 fencing 的 Redis SET NX PX 只能降低并发概率。
选主切换期间可能少一次也可能多一次,仍需执行账本和补跑。关键任务不要把“某个 Pod 应该一直活着”当可靠性证明。
13. 错过执行、补跑与时间窗口
进程从 02:50 停到 03:10,03:00 任务不会由 robfig/cron 自动补跑。启动时可读取执行账本,计算从最后成功窗口到当前应有的窗口,按上限补齐。必须限制追赶数量,避免停机十天后瞬间执行十天的昂贵任务。
任务函数最好接收显式 windowStart/windowEnd,而不是内部调用 time.Now() 猜本次处理范围。人工补跑使用相同入口和幂等键,并记录操作者、原因、代码版本。日窗口按业务时区生成,存储使用带时区的绝对时刻;不要用 24*time.Hour 推下一当地日。
补跑策略应区分可合并任务与不可合并任务:缓存刷新通常只跑最新窗口,财务日结可能逐窗执行,通知任务则可能因过期而直接终止。
14. 测试与诊断
把业务函数与 Cron 注册拆开。单元测试直接调用 run(ctx, window);计划测试使用 parser 的 Next 验证具体日期、月末、周末和 DST;生命周期集成测试用 channel 确认启动、重叠策略和 Stop 等待,不用长 Sleep。
func TestNextRunsUsesLocation(t *testing.T) {
loc, err := time.LoadLocation("Asia/Shanghai")
if err != nil {
t.Fatal(err)
}
from := time.Date(2026, 8, 31, 2, 59, 0, 0, loc)
runs, err := NextRuns("0 3 * * *", loc, from, 2)
if err != nil {
t.Fatal(err)
}
if runs[0].Hour() != 3 || runs[0].Day() != 31 {
t.Errorf("first run = %v", runs[0])
}
}
go test ./...
go test -race -count=20 ./...
go vet ./...
TZ=Europe/Berlin go test -run TestDST ./...
诊断“没跑”先看 scheduler 是否启动、Entry.Valid/Next、时区和表达式;再看 skipped、panic 和运行日志;最后查执行账本。诊断“跑多次”先确认副本数、默认重叠、重试/人工补跑和租约切换。
15. 指标、性能与容量
记录计划时间、实际开始、排队偏差、结束、状态、窗口和 run ID。指标包括触发数、成功/失败/skip、running、延迟、执行耗时、最后成功时间与错过窗口数。标签只用 job 名和稳定状态,不用 run ID。
Cron 自身开销通常很小,真正容量在 Job。大量 Entry 会增加排序和 timer 管理,但更常见的问题是同一分钟集中触发形成尖峰。通过错峰表达式、队列和下游并发限制平滑负载。Delay wrapper 造成的等待也占 goroutine,应监控而非当免费缓存。
任务读取大量数据时分页并固定上界,使用游标而非越来越慢的 offset;批次之间检查 context。内存缓存和连接池必须按多个同时运行任务的峰值估算。
16. 安全与生产部署边界
动态 cron 表达式是配置输入,必须限制长度、字段和最小频率,防止 * * * * * 类误配置压垮系统。任务参数需要验证和租户授权;日志不写 token、完整导出内容或 DSN。人工触发与补跑接口必须认证、授权、审计,并支持 dry-run 和速率限制。
Kubernetes CronJob 适合每次独立 Pod、资源隔离和平台级历史记录;设置 concurrencyPolicy、startingDeadlineSeconds、历史上限和 activeDeadlineSeconds,但它仍可能重复,业务仍需幂等。进程内 cron 适合轻量、与服务同生命周期且允许停机漏触发的任务。Asynq Scheduler 适合把触发与持久消费结合。
robfig/cron 的生产边界是:它可靠地计算进程在线期间的下一次时间,却不承诺持久执行、单副本、补跑或业务成功。 用 context 和 Stop 管理本地生命周期,用 wrapper 定义重入策略,用执行账本、幂等键、租约 fencing 或外部平台承担分布式语义,调度系统才可诊断、可恢复。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go Asynq 异步任务:Redis 队列、重试、定时与唯一任务
- 下一篇:Go Cobra CLI 实战:命令树、Flag、补全与可测试命令
- 延伸:Go 时间处理:time.Time、Duration、时区、Timer 与 Ticker
- 延伸:Go Redis 生产模式:缓存一致性、穿透击穿、锁与 Streams
- 延伸:Go 服务发现与配置中心:etcd、Consul、Nacos 的正确边界
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论