定时任务模块实现教程

2025-08-06T14:11:02+08:00 | 23分钟阅读 | 更新于 2025-08-06T14:11:02+08:00

@

学习目标

学完本章你应该能够:

  1. 讲清定时任务模块的分层设计:接口抽象(ScheduledTask / TaskManager)→ 实现层(cronTaskManager,基于 robfig/cron/v3)→ 生命周期集成(ScheduledTaskServer 适配 Kratos)。
  2. 说清 SkipIfStillRunningRecover 两个 cron 中间件的作用,以及它们"洋葱模型"式的包裹顺序与边界。
  3. 解释 RegisterStart 为什么分离,sync.Map + 读写锁如何保证并发安全。
  4. (商用核心) 指出单机 cron 在多副本部署下的重复执行隐患,并会用 Redis 分布式锁(TryLock + 看门狗续期)把"同任务同一时刻只有一个节点执行"落到代码。
  5. 用"基于状态重算 / 删除前先置位状态"让任务幂等,解释为什么"回收站清理"若不加锁会重复扣减存储。
  6. 处理时区、错过执行(backfill)、失败重试 + 告警 + 死信、优雅停机、Prometheus 可观测这五个生产级问题。

前置知识

  • Go 基础:goroutine、channel、sync.Map / sync.RWMutexcontext.Context、Redis 基本命令。
  • Kratos 框架基础:kratos.App 启动流程、transport.Server 接口。
  • robfig/cron/v3 基本用法(知道 CRON 表达式格式即可)。

本章你会动手做的事

  1. cron.New(cron.WithSeconds(), cron.WithChain(cron.SkipIfStillRunning(...), cron.Recover(...))) 起一个带防重叠 + panic 恢复的调度器。
  2. wrapFunc 加一层 Redis TryLock,模拟开两个进程,验证只有抢到锁的那个节点真正执行任务。
  3. 给存储校准任务加 ctx.Done() 监听与"按 user_id % N 分片",体验优雅退出与并行加速。

一、技术栈与中间件

定时任务模块底层用 robfig/cron/v3 做调度引擎,并通过 Kratos 的 transport.Server 接入应用生命周期。下表汇总用到的技术与中间件,并标注了商用化需要补强的地方:

技术 / 中间件用途说明
github.com/robfig/cron/v3Go 生态主流调度库,支持标准与秒级(6 字段)CRON 表达式,提供任务链(Chain)机制。本模块按表达式触发任务。
cron.WithSeconds()启用"秒级"CRON(6 字段:秒 分 时 日 月 周)。
cron.SkipIfStillRunning(logger)cron 内置 Job Wrapper:上一次未跑完则跳过本次触发,防任务重叠并发。
cron.Recover(logger)cron 内置中间件:捕获任务函数 panic,防单任务崩溃拖垮调度器。
cron.WithLocation(loc)(商用必加) 显式指定调度时区(如 Asia/Shanghai),避免各机器本地时区不一致导致执行时间漂移。
sync.Map并发安全 map,存 任务名 -> *taskWrapper
sync.RWMutex / sync.Mutex保护 taskWrapper.statusStartAll/StopAll 临界区。
context.Context任务函数签名 func(ctx context.Context) error 的取消与超时传递;StopAllcancel() 通知在跑任务退出。
transport.Server(Kratos)把调度器适配成"服务",复用框架优雅启停。
github.com/google/wire编译期依赖注入,组装 TaskManagerScheduledTaskServer
log/slog结构化日志(开始 / 完成 / 失败 / 耗时)。
data/lock 分布式锁(项目已有)(商用必加) 基于 Redis SET NX + Lua 的 TryLock,带 TTL 自动释放与后台续期(renewLoop)。用于多实例下保证"同任务单点执行"。
Prometheus 指标(需补)(商用必加) 记录任务耗时、成功率、最近成功时间戳,供 Grafana / AlertManager 监控告警。

二、实现思路流程(总体)

定时任务模块遵循"分层抽象 + 生命周期托管 + 单点执行保护"的设计,从底层调度器到业务任务逐层封装:

  1. cron 调度器初始化(SkipIfStillRunning + Recover + 时区)

    • NewCronTaskManager 通过 cron.New(cron.WithSeconds(), cron.WithChain(cron.SkipIfStillRunning(...), cron.Recover(...))) 创建调度器。修正点:商用版需追加 cron.WithLocation(time.LoadLocation("Asia/Shanghai")) 固定时区。
  2. 任务定义(ScheduledTask 结构)

    • 定义 ScheduledTask 接口(Name/Spec/Func/Status/Start/Stop/Restart)与 TaskFuncfunc(ctx context.Context) error),NewTask 便捷构造。
  3. 任务注册(sync.Map 管理)

    • cronTaskManagersync.Map 维护 任务名 -> *taskWrapperRegister 去重(同名返回 ErrTaskAlreadyExists),但不立即加入 cron,需等 Start/StartAll`。
  4. 单点执行保护(分布式锁,商用必加)

    • 修正点:在多实例部署下,每个 pod 都会各自跑一遍 cron。必须在任务被触发时先用 Redis TryLock("cloud-disk:scheduler🔒{name}") 抢占,抢不到就跳过;否则回收站清理会被多节点并发执行,导致存储被重复扣减、事件被重复发布。
  5. 任务启动 / 优雅停机

    • StartAll 通过 cron.AddFunc 注册并 cron.Start()wrapFunc 统一做"状态流转 + 耗时统计 + 锁/可观测/告警"。StopAllcron.Stop() 并在 ctxcancel(),最多等 10 秒让在跑任务退出。
  6. ScheduledTaskServer 集成 Kratos 生命周期

    • ScheduledTaskServer 实现 transport.ServerStart→StartAllStop→StopAll,随 app.Run() 启停,无需业务代码手动管理。
  7. 具体业务任务

    • newStorageCalibrateTask(每天 3:00):全量重算用户已用空间(对账纠偏)。
    • newRecycleCleanTask(每天 3:00):删过期回收站记录 + 物理删文件 + 扣减存储。
    • newChunkCleanupTask(每天 4:00):删中断上传残留的孤儿分片目录。
    • 缺失能力(商用补全):过期分享链接目前只在"访问时"判过期(ErrShareExpired),没有定时任务主动失效并清理,应补 share-clean 任务。

下面这张时间轴图把三个业务任务的 CRON 时间表画出来,注意 3 点两个任务并发、4 点错峰:

flowchart LR
    T03[每天 03:00] --> A[存储校准
遍历全量用户] T03 --> B[回收站清理
删过期记录+扣存储] T04[每天 04:00] --> C[孤儿分片清理
删残留 chunks]

整体调用链:main.go newTaskManagerRegister 记账 → ScheduledTaskServer.Start(Kratos 启动)→ StartAllcron.Start → 到点触发 wrapFunc(含分布式锁抢占)→ 业务 TaskFunc

flowchart TD
    M[main.go] -->|newTaskManager 注册| TM[cronTaskManager]
    M -->|注入| STS[ScheduledTaskServer]
    STS -->|实现 transport.Server| K[Kratos App 生命周期]
    K -->|Start 时| SA[StartAll]
    SA -->|AddFunc 注册| C[cron 调度器
WithLocation 时区] C -->|到点触发| WF[wrapFunc
分布式锁 + 状态 + 日志 + 指标] WF -->|抢占成功| TF[业务 TaskFunc] WF -->|未抢到锁 跳过| SK[其他节点已在跑]

三、面试常问知识点与难点

1. CRON 表达式语法

本模块用 cron.WithSeconds() 启用 6 字段:秒 分 时 日 月 周0 0 3 * * * = 每天 03:00:00;0 0 4 * * * = 每天 04:00:00;*/30 * * * * * = 每 30 秒。* 任意值,*/n 每 n 单位,1-5 区间。注意秒级表达式与 5 字段(不含秒)的区别,面试常考。

2. SkipIfStillRunning 防止任务重叠

若任务执行时间超过触发间隔,默认 cron 会再起一个实例造成并发。SkipIfStillRunning 持一把内部锁,检测到上次仍在跑就跳过本次并记日志。对"遍历全表"类任务很重要,可避免并发写冲突。但注意:它只防"同一进程内"的重叠,不防"多进程/多实例"——跨节点重复执行要靠分布式锁(见 3.4)。

3. Recover 中间件 panic 恢复

Go 中 goroutine panic 未 recover 会让整个进程崩溃。cron.Recover(logger) 在每个任务外层包 defer recover(),把 panic 转日志,确保单任务崩溃不拖垮调度器。这是生产定时任务必备兜底。

⚠️ 洋葱顺序的边界WithChain(SkipIfStillRunning, Recover)SkipIfStillRunning 在最外层、Recover 在内层。意味着"是否还在跑"的判断包在 panic 恢复之外——若任务 panic,本次仍会被记为"运行中"直到 Recover 处理完,下一轮才会正常判断。两者职责不重叠,顺序不可颠倒。

4. 多实例重复执行(分布式调度的真实隐患)

单机 cron 在 K8s / 多副本部署时,每个 pod 都会各自起一个 cron 实例,于是三个任务在 N 个节点上各跑一遍:

  • 孤儿分片清理:文件删除幂等,影响较小,但会重复打日志、重复发布事件;
  • 存储校准RecalibrateStorage 内部用一把全局 calibrateLockKey 锁,事实上把多节点的校准串行化了(保了命但成了性能瓶颈,且不是为"去重"而设计);
  • 回收站清理CleanExpired 物理删文件并调用 SubUsedStorage 扣存储,没有任何跨节点互斥——两个节点同时跑,同一批过期记录会被处理两次,存储被重复扣减,用户配额出现负向偏差。

项目里其实已经有 Redis 分布式锁实现(data/lock/redis.goTryLock,带 TTL 自动释放 + renewLoop 后台续期),但调度器并没有用它。正确做法:任务被触发时先 TryLock("cloud-disk:scheduler🔒{name}"),抢到才执行,抢不到直接跳过。分布式方案也可换成 XXL-JOB / Elastic-Job / asynq / K8s CronJob(见 3.2)。

5. 任务幂等性

定时任务可能因重试、重启、多节点被多次执行,逻辑必须幂等。本项目中:

  • 存储校准是"重算覆盖"(UPDATE SET used = 计算值 + delta),结果一致;
  • 孤儿分片清理是"删文件",幂等;
  • 回收站清理是隐患:它"查过期 → 物理删 + 扣存储 → 删库记录",在没锁的情况下重复执行会重复扣存储。加固方式:① 加分布式锁(主);② 删除前先把记录标记 status=deleted(乐观锁 WHERE status=...),即使两节点并发也只有一方生效。

⚠️ 新手必踩:别在定时任务里写"累加 / 插入"类非幂等操作。比如"给每个用户加 100 积分"放进定时任务,一旦重试或重跑就多加。正确做法是"基于状态重算"或"先 SELECT 判断是否已处理"。

6. CRON 时区问题

robfig/cron/v3 默认用机器本地时区。若容器镜像默认 UTC(FROM alpine)而代码以为本地是 Asia/Shanghai0 0 3 * * * 的实际触发时刻就会错位。最稳妥:代码 cron.WithLocation(time.LoadLocation("Asia/Shanghai")) + 容器设 TZ 环境变量 / 挂时区文件,双管齐下。

7. 任务失败重试、告警与死信

本模块 wrapFunc 仅把状态置 TaskStatusFailed 并打日志,不重试、不告警、无死信。生产任务(尤其外部调用、批量删除)应:

  • 区分"瞬时失败可重试"(网络抖动)与"业务失败不可重试"(数据错误);
  • 瞬时失败走指数退避重试(retry/backoff);
  • 重试耗尽写死信表task_dead_letter)并触发告警(Webhook / 钉钉 / 邮件);
  • 单任务失败不能中断整体(如存储校准逐用户容错)。

8. 错过执行(backfill)——宕机期间的任务补跑

robfig/cron 默认不在进程宕机后补跑错过的调度:若服务在 03:00 宕机,当天的存储校准/回收站清理就直接跳过,过期回收站多占一天空间、存储偏差多留一天。商用需要"错过即补":

  • 用 Redis 记录每个任务的 last_success_atcloud-disk:scheduler:last:{name});
  • 进程启动 StartAll 时,若 now - last_success > 调度周期 且错过了窗口,立即补跑一次;
  • 或对强一致任务用 K8s CronJob(自带 startingDeadlineSeconds backfill)。

9. Kratos transport.Server 生命周期集成

transport.Server 只有 Start(ctx) / Stop(ctx)ScheduledTaskServer 实现该接口,把定时任务当成一种"服务器"注入 kratos.App,复用优雅启停、信号处理、依赖注入。这是一个"适配器模式":把任意可启停组件适配成 transport.Server。

10. 可观测性(metrics)

slog 打印耗时不够。商用应暴露 Prometheus 指标:

  • scheduler_task_duration_seconds{name}(耗时直方图);
  • scheduler_task_last_success_ts{name}(最近成功时间戳,用于"多久没跑"告警);
  • scheduler_task_failures_total{name}(失败计数,连续失败触发 AlertManager)。

四、生产化加固(商用审查重点)

本章把"多实例重复、幂等、失败重试、时区/错过执行、可观测、优雅停机"落到可执行的加固项。关键修正直接进第五节代码(标「修正点」),不另列改造清单。

1. 分布式锁保证多实例单点执行

cronTaskManager 注入 lock.LockerwrapFunc 触发时先 TryLock("cloud-disk:scheduler🔒{name}", TTL=任务预估时长);抢到才执行,抢不到 log.Info 跳过。TTL 用后台 renewLoop 续期,防任务跑太久锁过期被别的节点抢走(项目 redisLock 已支持)。

2. 幂等加固(删除前先置位状态)

CleanExpired 改为"先 UPDATE recycle_bins SET status=deleted WHERE expire_at<? AND status=active 拿到受影响行 → 仅对这批做物理删 + 扣存储"。即使两个节点同时进入,乐观锁保证只有一方真正处理,存储不会被重复扣减。存储校准本就是覆盖写,天然幂等。

3. 失败重试 + 告警 + 死信

wrapFunc 内对任务返回的 error 分类:瞬时错误指数退避重试(最多 3 次);重试耗尽或业务错误写死信表并调 alert(ctx, name, err)。监控看见 scheduler_task_failures_total 飙升即告警。

4. 时区与错过执行(backfill)

NewCronTaskManagerloc *time.LocationStartAll 时读取 last_success,对"错过窗口"的任务立即补跑。Recover/锁异常时不影响补跑判定。

5. 任务分片并行 + 限流

存储校准从"单 goroutine 串行遍历全量用户"改为按 user_id % N 分片、N 个 worker 并行;对下游 DB / 对象存储调用用 semaphore.Weightedrate.Limiter 限流,防打垮依赖。循环中监听 ctx.Done() 支持优雅退出。

6. 可观测(Prometheus + 告警)

wrapFunc 埋点:记录耗时直方图、成功/失败计数、最近成功时间戳。配合 Grafana 面板与 AlertManager:任务超过阈值未成功(如超 26h 没跑)即告警。

7. 优雅停机与断点续跑

StopAll 已是 cron.Stop() + cancel(ctx) + 10s 等待。长任务应在循环里查 ctx.Done() 提前退出;对超大扫描任务,把游标(last_user_id)持久化,重启后从断点续跑,避免每次从头。


五、详细实现流程与代码解析

5.1 Cron 调度器初始化(SkipIfStillRunning + Recover + 时区 + 锁/指标注入)

实现思路

NewCronTaskManager 创建 cron.Cron 并配置:

  • cron.WithSeconds():秒级表达式;
  • cron.WithChain(SkipIfStillRunning, Recover):先判重叠、再兜底 panic;
  • 修正点:追加 cron.WithLocation(loc) 固定时区;
  • 修正点:持有 locker(分布式锁)与 metrics(指标),供 wrapFunc 使用。

洋葱模型:越外层越先执行。SkipIfStillRunning 最外层,其次 Recover,再 wrapFunc(锁/状态/指标),最里是业务 TaskFunc

flowchart TD
    A[cron 触发] --> B[SkipIfStillRunning
上次没跑完就跳过] B --> C[Recover
兜住 panic] C --> D[wrapFunc
分布式锁 + 状态 + 指标 + 日志] D -->|抢到锁| E[业务 TaskFunc] D -->|未抢到锁| S[跳过 记日志]

关键代码(internal/data/scheduler/cron_scheduler.go,商用修正版)

package scheduler

import (
	"context"
	"log/slog"
	"sync"
	"time"

	"github.com/go-kratos/kratos/v3/log"
	"github.com/google/wire"
	"github.com/robfig/cron/v3"
	"your-module/internal/data/lock" // 复用项目已有的 Redis 分布式锁
)

// TaskMetrics 封装定时任务的 Prometheus 指标(商用必加)。
type TaskMetrics struct {
	Duration *prometheus.HistogramVec // scheduler_task_duration_seconds
	Failures *prometheus.CounterVec   // scheduler_task_failures_total
	LastOK   *prometheus.GaugeVec     // scheduler_task_last_success_ts
}

// Option 修正点:用选项模式注入时区、分布式锁、指标,避免改构造函数签名破坏 wire。
type Option func(*cronTaskManager)

func WithLocation(loc *time.Location) Option {
	return func(m *cronTaskManager) { m.location = loc }
}
func WithLocker(l lock.Locker) Option {
	return func(m *cronTaskManager) { m.locker = l }
}
func WithMetrics(mt *TaskMetrics) Option {
	return func(m *cronTaskManager) { m.metrics = mt }
}

type cronTaskManager struct {
	cron     *cron.Cron
	tasks    sync.Map
	logger   *slog.Logger
	ctx      context.Context
	cancel   context.CancelFunc
	mu       sync.Mutex
	location *time.Location  // 修正点:显式时区
	locker   lock.Locker     // 修正点:分布式锁,保证多实例单点执行
	metrics  *TaskMetrics    // 修正点:可观测
}

// NewCronTaskManager 修正点:默认时区改为 Asia/Shanghai,并支持注入锁与指标。
func NewCronTaskManager(opts ...Option) TaskManager {
	ctx, cancel := context.WithCancel(context.Background())
	loc, _ := time.LoadLocation("Asia/Shanghai") // 修正点:固定时区,不依赖机器本地

	c := cron.New(
		cron.WithSeconds(),
		cron.WithLocation(loc), // 修正点:时区显式指定
		cron.WithChain(
			cron.SkipIfStillRunning(cron.DiscardLogger),
			cron.Recover(cron.DefaultLogger),
		),
	)

	cm := &cronTaskManager{
		cron:     c,
		logger:   log.Default(),
		ctx:      ctx,
		cancel:   cancel,
		location: loc,
	}
	for _, o := range opts {
		o(cm)
	}
	return cm
}

要点:

  • 修正点cron.WithLocation(loc) 让所有节点"在同一时刻"触发,消除容器时区漂移。
  • 修正点lockermetrics 通过 Option 注入,既满足 wire 兼容,又把商用能力挂到 wrapFunc

5.2 任务定义与注册(ScheduledTask + TaskManager)

实现思路

scheduler.go 定义 ScheduledTask 接口、TaskManager 接口、TaskFunc 签名、TaskStatus 枚举、NewTask + simpleTasktaskWrappersimpleTask 上增加 cron.EntryID 与带锁 statusRegistersync.Map 去重记账,不立即调度。

关键代码(internal/data/scheduler/scheduler.go

// TaskStatus 表示定时任务的当前状态。
type TaskStatus int32

const (
	TaskStatusStopped TaskStatus = 0 // 已注册未运行
	TaskStatusRunning TaskStatus = 1 // 已调度且活动
	TaskStatusFailed  TaskStatus = 2 // 上次执行出错
)

// TaskFunc 是定时任务执行函数的统一签名。
type TaskFunc func(ctx context.Context) error

// ScheduledTask 定义可被调度的任务的接口。
type ScheduledTask interface {
	Name() string
	Spec() string
	Func() TaskFunc
	Status() TaskStatus
	Start() error
	Stop() error
	Restart() error
}

// TaskManager 任务管理器接口(注册/启停/查询)。
type TaskManager interface {
	Register(task ScheduledTask) error
	Unregister(name string) error
	Start(name string) error
	Stop(name string) error
	Restart(name string) error
	GetTask(name string) (ScheduledTask, bool)
	ListTasks() []ScheduledTask
	StartAll() error
	StopAll() error
}

var (
	ErrTaskNotFound      = errors.New("task not found")
	ErrTaskAlreadyExists = errors.New("task already registered")
	ErrTaskNotRunning    = errors.New("task is not running")
	ErrTaskAlreadyActive = errors.New("task is already active")
)

// NewTask 便捷构造器。
func NewTask(name, spec string, fn TaskFunc) ScheduledTask {
	return &simpleTask{name: name, spec: spec, fn: fn}
}

cron_scheduler.gotaskWrapperstatus 用读写锁保护(真正并发安全版),simpleTask 的无锁版本用于临时任务:

type taskWrapper struct {
	name   string
	spec   string
	fn     TaskFunc
	status TaskStatus
	cronID cron.EntryID
	mu     sync.RWMutex
}

func (t *taskWrapper) Status() TaskStatus {
	t.mu.RLock()
	defer t.mu.RUnlock()
	return t.status
}
func (t *taskWrapper) setStatus(s TaskStatus) {
	t.mu.Lock()
	defer t.mu.Unlock()
	t.status = s
}

// Register 记账去重,不调用 cron.AddFunc。
func (m *cronTaskManager) Register(task ScheduledTask) error {
	name := task.Name()
	if _, ok := m.tasks.Load(name); ok {
		return ErrTaskAlreadyExists
	}
	tw := &taskWrapper{name: name, spec: task.Spec(), fn: task.Func(), status: TaskStatusStopped}
	m.tasks.Store(name, tw)
	m.logger.Info("scheduler: task registered", "name", name, "spec", tw.spec)
	return nil
}

要点:Register 只记账,Start/StartAllAddFunc,保证"先批量注册、再统一启动"的状态一致。


5.3 ScheduledTaskServer 生命周期与优雅停机

实现思路

ScheduledTaskServer 实现 Kratos transport.ServerStart→StartAllStop→StopAllStopAllcron.Stop()(停收新触发、允许在跑任务跑完),再 m.cancel() 取消任务 ctx,最多等 10 秒。

sequenceDiagram
    participant App as Kratos App
    participant STS as ScheduledTaskServer
    participant M as cronTaskManager
    participant C as cron 调度器
    App->>STS: Start(ctx)
    STS->>M: StartAll()
    M->>C: AddFunc + cron.Start() + backfill 补跑
    Note over C: 到点触发 wrapFunc(含分布式锁)
    App->>STS: Stop(ctx) 退出信号
    STS->>M: StopAll()
    M->>C: cron.Stop() + cancel(ctx)
    C-->>M: 在跑任务退出后 Done

关键代码(优雅停机 + 双启动守卫)

// StartAll 启动所有任务,并做 backfill 补跑(修正点)。
func (m *cronTaskManager) StartAll() error {
	m.mu.Lock()
	defer m.mu.Unlock()
	if m.cron.Running() { // 修正点:守卫,避免重复 Start 引发 cron 报错
		return nil
	}

	var lastErr error
	m.tasks.Range(func(key, value interface{}) bool {
		tw := value.(*taskWrapper)
		if tw.cronID == 0 {
			id, err := m.cron.AddFunc(tw.spec, m.wrapFunc(tw))
			if err != nil {
				m.logger.Error("scheduler: failed to start task", "name", tw.name, "err", err)
				lastErr = err
				return true
			}
			tw.cronID = id
			tw.setStatus(TaskStatusRunning)
		}
		return true
	})

	m.cron.Start()

	// 修正点:错过执行补跑——若上次成功距今已超过一个周期,立即跑一次
	m.tasks.Range(func(key, value interface{}) bool {
		tw := value.(*taskWrapper)
		if m.shouldBackfill(tw.name, tw.spec) {
			go func(t *taskWrapper) {
				m.logger.Info("scheduler: backfill run", "name", t.name)
				t.fn(m.ctx)
			}(tw)
		}
		return true
	})

	m.logger.Info("scheduler: all tasks started")
	return lastErr
}

// StopAll 优雅停机。
func (m *cronTaskManager) StopAll() error {
	m.mu.Lock()
	defer m.mu.Unlock()

	ctx := m.cron.Stop() // 停止接收新触发,等待在跑任务完成
	m.tasks.Range(func(key, value interface{}) bool {
		tw := value.(*taskWrapper)
		tw.cronID = 0
		tw.setStatus(TaskStatusStopped)
		return true
	})
	m.cancel() // 通知在跑任务主动退出

	select {
	case <-ctx.Done():
		m.logger.Info("scheduler: all tasks stopped gracefully")
	case <-time.After(10 * time.Second):
		m.logger.Warn("scheduler: timed out waiting for tasks to complete")
	}
	return nil
}

要点:修正点if m.cron.Running() 守卫防止 StartAll 被重复调用(Kratos 某些重启路径)导致 cron 报错;shouldBackfill 实现"宕机期间错过的任务,启动即补跑"。


5.4 存储校准定时任务(每天 3:00,含 ctx 检查与分片)

实现思路

云盘日常增删文件更新 used_storage,异常会导致偏差。存储校准每天 3:00 全量重算纠偏。RecalibrateStorage 内部已是"按实际活跃文件大小重算 + delta 覆盖",天然幂等;但它是串行遍历全量用户,且循环未查 ctx.Done()

修正点:① 循环里监听 ctx.Done() 优雅退出;② 按 user_id % N 分片并行;③ 限流下游 DB。多实例下由 wrapFunc 的分布式锁保证只有一个节点跑。

关键代码(cmd/server/main.go,修正版)

// newStorageCalibrateTask 每天 3:00 存储校准。
func newStorageCalibrateTask(userUC *biz.UserUsecase) scheduler.ScheduledTask {
	return scheduler.NewTask(
		"storage-calibrate",
		"0 0 3 * * *",
		func(ctx context.Context) error {
			slog.Info("storage calibrate task running")
			ids, err := userUC.ListAllUserIDs(ctx)
			if err != nil {
				return err
			}

			// 修正点:按 user_id % N 分片并行,semaphore 限制并发,循环监听 ctx.Done()
			const workers = 8
			sem := semaphore.NewWeighted(workers)
			var wg sync.WaitGroup
			for _, id := range ids {
				// 修正点:优雅退出——收到取消立即结束
				if ctx.Err() != nil {
					break
				}
				if err := sem.Acquire(ctx, 1); err != nil {
					break
				}
				wg.Add(1)
				go func(uid uint64) {
					defer wg.Done()
					defer sem.Release(1)
					if err := userUC.RecalibrateStorage(ctx, uid); err != nil {
						slog.Error("storage calibrate: failed", "user_id", uid, "err", err)
					}
				}(id)
			}
			wg.Wait()
			slog.Info("storage calibrate task completed", "users", len(ids))
			return nil
		},
	)
}

要点:RecalibrateStorage 用全局 calibrateLockKey 串行化同一时刻的校准(含跨节点),但性能受限;生产应把锁粒度细化到"单用户"并在任务外层用节点级分布式锁,既去重又并行。


5.5 回收站过期清理定时任务(幂等 + 锁)

实现思路

用户删除文件进入回收站,保留 7 天可恢复,超期需物理清理(删 DB 记录 + 底层文件块 + 扣存储)。这是多实例重复执行的高危任务:无跨节点互斥时会重复 SubUsedStorage,造成配额负向偏差。

修正点:① wrapFunc 分布式锁保证单点执行(主);② CleanExpired 内部"先乐观锁置位 status=deleted 再处理",即便两节点并发也只一方生效。

关键代码(internal/biz/recycle.go,CleanExpired 幂等加固版)

// CleanExpired 删除已过期回收站项目(定时任务调用)。
// 修正点:先乐观锁置位,再处理,保证即使无分布式锁也不会重复扣存储。
func (uc *RecycleUsecase) CleanExpired(ctx context.Context) error {
	// 1) 先批量把"已过期且仍 active"的记录置为 deleted(行级乐观锁,幂等)
	affected, err := uc.recycleRepo.MarkExpiredDeleted(ctx, time.Now())
	if err != nil {
		return err
	}
	if affected == 0 {
		return nil // 没有需要处理,直接返回(幂等)
	}
	// 2) 只处理已被本节点置位的这批(status=deleted 且未物理删)
	expired, err := uc.recycleRepo.ListMarkedDeleted(ctx)
	if err != nil {
		return err
	}

	deletedByUser := make(map[uint64][]uint64)
	for _, item := range expired {
		ids := uc.permanentlyDeleteItem(ctx, item.UserID, item) // 物理删 + SubUsedStorage
		if len(ids) > 0 {
			deletedByUser[item.UserID] = append(deletedByUser[item.UserID], ids...)
		}
		if err := uc.recycleRepo.Delete(ctx, item.ID); err != nil {
			log.Error("recycle: failed to delete item", "id", item.ID, "err", err)
		}
	}

	if uc.eventPublisher != nil {
		for userID, fileIDs := range deletedByUser {
			_ = uc.eventPublisher.Publish(ctx, EventFileDeleted, &FileChangedPayload{
				UserID: userID, FileIDs: fileIDs, Action: "deleted",
			})
		}
	}
	return nil
}

要点:修正点MarkExpiredDeletedUPDATE ... SET status=deleted WHERE expire_at<? AND status=active 的受影响行数代表"本节点认领的工作量",天然规避重复扣存储;配合 wrapFunc 的 Redis 锁,双保险。


5.6 孤儿分片目录清理定时任务

实现思路

大文件分片上传,若上传中崩溃/客户端断开未取消,会残留 uploadDir/chunks/<uploadID>/ 孤儿目录。每天 4:00(错峰)调 CleanupOrphanChunks(ctx, 24h) 删修改时间超 24h 的目录——24h 阈值避免误删进行中的上传。文件删除幂等,重复执行安全;但仍应走分布式锁避免多节点无谓重复扫描与重复事件。

关键代码(cmd/server/main.go

func newChunkCleanupTask(storage biz.Storage) scheduler.ScheduledTask {
	return scheduler.NewTask(
		"chunk-cleanup",
		"0 0 4 * * *",
		func(ctx context.Context) error {
			slog.Info("chunk cleanup task running")
			removed, err := storage.CleanupOrphanChunks(ctx, 24*time.Hour)
			if err != nil {
				slog.Error("chunk cleanup: failed", "err", err)
				return err
			}
			slog.Info("chunk cleanup task completed", "removed_count", len(removed))
			return nil
		},
	)
}

要点:24h 是业务安全边界;返回值 removed 用于审计,若某天 removed_count 异常飙升,可能说明上传通道异常。


5.7 用 Redis 分布式锁保证单点执行(wrapFunc 集成锁 + 可观测 + 告警/死信 + backfill)

这是把"商用审查重点"一次性落地的核心。wrapFuncRecover 之内、TaskFunc 之外,包一层:抢锁 → 记录指标 → 执行(含重试)→ 释放锁 → 写 last_success / 失败告警 / 死信。

flowchart TD
    T[cron 触发 wrapFunc] --> L[TryLock cloud-disk:scheduler🔒name]
    L -->|未抢到| SK[跳过 记日志 返回]
    L -->|抢到| R[renewLoop 后台续期]
    R --> RUN[执行业务 TaskFunc]
    RUN --> OK{成功?}
    OK -->|是| W1[写 last_success + 指标]
    OK -->|否| RETRY{瞬时错误且可重试?}
    RETRY -->|是| RUN2[指数退避重试 最多3次]
    RETRY -->|否| DL[写死信表 + 告警]
    W1 --> U[Unlock]
    DL --> U
    SK --> X[结束]
    U --> X

关键代码(internal/data/scheduler/cron_scheduler.go,wrapFunc 商用版)

// wrapFunc 包装任务:分布式锁 + 状态 + 指标 + 重试 + 告警 + 死信。
func (m *cronTaskManager) wrapFunc(tw *taskWrapper) cron.FuncJob {
	return func() {
		tw.setStatus(TaskStatusRunning)

		// 修正点:多实例单点执行——抢不到锁直接跳过
		lockKey := "cloud-disk:scheduler🔒" + tw.name
		if m.locker != nil {
			got, err := m.locker.TryLock(m.ctx, lockKey,
				lock.WithTTL(30*time.Minute), lock.WithRenew()) // 带后台续期
			if err != nil || !got {
				m.logger.Info("scheduler: another node holds lock, skip", "name", tw.name)
				return
			}
			defer func() { _ = m.locker.Unlock(m.ctx, lockKey) }()
		}

		start := time.Now()
		err := m.runWithRetry(tw) // 修正点:带重试
		duration := time.Since(start)

		// 修正点:可观测埋点
		if m.metrics != nil {
			m.metrics.Duration.WithLabelValues(tw.name).Observe(duration.Seconds())
			if err == nil {
				m.metrics.LastOK.WithLabelValues(tw.name).Set(float64(time.Now().Unix()))
			} else {
				m.metrics.Failures.WithLabelValues(tw.name).Inc()
			}
		}

		if err != nil {
			tw.setStatus(TaskStatusFailed)
			m.logger.Error("scheduler: task failed", "name", tw.name, "duration", duration, "err", err)
			m.alert(tw.name, err)   // 修正点:告警
			m.writeDeadLetter(tw.name, err) // 修正点:死信
			return
		}
		tw.setStatus(TaskStatusRunning)
		m.logger.Info("scheduler: task completed", "name", tw.name, "duration", duration)
	}
}

// runWithRetry 瞬时错误指数退避重试,业务错误直接返回。
func (m *cronTaskManager) runWithRetry(tw *taskWrapper) error {
	const maxRetry = 3
	var err error
	for i := 0; i < maxRetry; i++ {
		if m.ctx.Err() != nil {
			return m.ctx.Err()
		}
		err = tw.fn(m.ctx)
		if err == nil {
			return nil
		}
		if !isTransient(err) { // 业务错误不重试
			return err
		}
		backoff := time.Duration(1<<uint(i)) * time.Second // 1s,2s,4s
		select {
		case <-time.After(backoff):
		case <-m.ctx.Done():
			return m.ctx.Err()
		}
	}
	return err
}

要点:

  • 修正点TryLock(... WithRenew()) 复用项目 redisLock 的续期能力,任务跑久也不怕锁过期被别节点抢走。
  • 修正点runWithRetry 区分瞬时/业务错误,避免对"数据错误"无效重试。
  • 修正点:失败写死信 + 告警,使"单任务失败不丢、可复盘"。
  • 修正点last_success / 指标让"任务多久没跑"可被监控告警。

5.8 过期分享清理(缺失能力补全)

现状:分享链接有 ExpireAt,但只在访问时ErrShareExpiredshare.go:225),没有定时任务主动把过期分享置为失效、回收 DB 与索引。商用应补 share-clean 任务,逻辑与回收站清理同理(乐观锁置位 + 幂等 + 分布式锁)。

// newShareCleanTask 每天 2:00 清理过期分享(缺失能力补全)。
func newShareCleanTask(shareUC *biz.ShareUsecase) scheduler.ScheduledTask {
	return scheduler.NewTask(
		"share-clean",
		"0 0 2 * * *",
		func(ctx context.Context) error {
			slog.Info("share clean task running")
			if err := shareUC.CleanExpired(ctx); err != nil {
				slog.Error("share clean: failed", "err", err)
				return err
			}
			slog.Info("share clean task completed")
			return nil
		},
	)
}

要点:过期分享不主动清理只会造成"分享表无限膨胀 + 用户看到已失效的分享",属于商用版应补的定时任务。其实现同样遵循"幂等 + 分布式锁"原则。


自测题与动手练习

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

  1. SkipIfStillRunningRecover 各自解决什么?在 WithChain 里包裹顺序如何影响行为?它们能防"多实例重复执行"吗?
  2. 多副本部署下,本模块原版的三个任务各节点都会跑。哪个任务因此会重复扣减存储?项目已有但未接入的什么能力可以解决它?请画出抢锁流程图。
  3. 回收站清理要做到"即使没锁也不重复扣存储",业务层该怎么改?(提示:乐观锁置位 status=deleted
  4. robfig/cron 在进程宕机期间错过的调度会补跑吗?商用如何做 backfill?请用一句话说明。
  5. NewCronTaskManager() 原版用默认本地时区有什么风险?容器镜像默认时区常是什么?正确做法?
  6. wrapFunc 原版只在失败记日志。生产应加哪三件事(重试 / 告警 / 死信),并区分哪两类错误?
  7. 存储校准任务要支持"优雅退出 + 并行加速",循环里要监听什么?按什么字段分片?用什么限制下游并发?
  8. 为什么存储校准本身是幂等的,而回收站清理原版不是?分别说明原因。

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

  1. cron.New(cron.WithSeconds(), cron.WithChain(cron.SkipIfStillRunning(cron.DiscardLogger), cron.Recover(cron.DefaultLogger))),写个 1 秒一次、故意 time.Sleep(3s) 的任务,观察是否被跳过。
  2. wrapFunc 加一层 Redis TryLock("cloud-disk:scheduler🔒demo"),开两个进程跑同一任务,验证只有抢到锁的节点真正执行、另一个记日志跳过。
  3. 把存储校准从"串行遍历"改成"按 user_id % 8 分 8 个 worker + semaphore",循环里加 if ctx.Err() != nil { break },跑起来后 StopAll 验证能优雅退出而非强杀。
  4. 给回收站清理加 MarkExpiredDeleted(乐观锁 UPDATE ... SET status=deleted WHERE expire_at<? AND status=active),用两个 goroutine 并发调,验证只有一方认领、存储不被重复扣减。

本章小结

  • 分层与生命周期ScheduledTask / TaskManager 接口屏蔽底层,cronTaskManager 基于 robfig/cron/v3 + sync.Map 落地,ScheduledTaskServer 用适配器模式把调度器塞进 Kratos 生命周期;SkipIfStillRunning + Recover + wrapFunc 保证不重叠、不崩溃、可观测。
  • 多实例单点执行(核心修正):原版无分布式锁,多副本会重复执行,回收站清理会重复扣减存储。商用修正为 wrapFunc 触发时 TryLock("cloud-disk:scheduler🔒{name}") + 后台续期,抢不到即跳过;业务层再用乐观锁置位 status=deleted 做幂等双保险。
  • 时区与错过执行cron.WithLocation(Asia/Shanghai) 固定时区,消除容器时区漂移;StartAll 做 backfill,宕机期间错过的任务启动即补跑。
  • 失败与可观测wrapFunc 加指数退避重试(区分瞬时/业务错误)、失败告警、死信表,并埋 Prometheus 指标(耗时 / 最近成功时间戳 / 失败计数),使"任务多久没跑、跑得快慢、是否连续失败"可被监控。
  • 优雅停机与并行StopAllcron.Stop() + cancel(ctx) + 10s 等待;长任务循环监听 ctx.Done(),存储校准按 user_id % N 分片并行 + semaphore 限流。
  • 缺失能力:过期分享仅访问时判过期,应补 share-clean 定时任务(同样遵循幂等 + 分布式锁)。
  • 下一篇可深入"分布式任务调度框架选型"(XXL-JOB / Elastic-Job / asynq / K8s CronJob),把本模块的抽象真正扩展到多机集群调度。
About Me

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

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

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

目标

学AI,加油!加油!