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/ShanghaiEurope/BerlinCST 有歧义,固定 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 的对象;两者返回稳定的 EntryIDEntry(id) 返回快照,可查看 PrevNextValidEntries() 用于诊断;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、资源隔离和平台级历史记录;设置 concurrencyPolicystartingDeadlineSeconds、历史上限和 activeDeadlineSeconds,但它仍可能重复,业务仍需幂等。进程内 cron 适合轻量、与服务同生命周期且允许停机漏触发的任务。Asynq Scheduler 适合把触发与持久消费结合。

robfig/cron 的生产边界是:它可靠地计算进程在线期间的下一次时间,却不承诺持久执行、单副本、补跑或业务成功。 用 context 和 Stop 管理本地生命周期,用 wrapper 定义重入策略,用执行账本、幂等键、租约 fencing 或外部平台承担分布式语义,调度系统才可诊断、可恢复。


系列导航与关联阅读

官方资料

本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。