学习目标
学完本章你应该能够:
- 讲清定时任务模块的分层设计:接口抽象(
ScheduledTask/TaskManager)→ 实现层(cronTaskManager,基于robfig/cron/v3)→ 生命周期集成(ScheduledTaskServer适配 Kratos)。 - 说清
SkipIfStillRunning与Recover两个 cron 中间件的作用,以及它们"洋葱模型"式的包裹顺序与边界。 - 解释
Register与Start为什么分离,sync.Map+ 读写锁如何保证并发安全。 - (商用核心) 指出单机 cron 在多副本部署下的重复执行隐患,并会用 Redis 分布式锁(
TryLock+ 看门狗续期)把"同任务同一时刻只有一个节点执行"落到代码。 - 用"基于状态重算 / 删除前先置位状态"让任务幂等,解释为什么"回收站清理"若不加锁会重复扣减存储。
- 处理时区、错过执行(backfill)、失败重试 + 告警 + 死信、优雅停机、Prometheus 可观测这五个生产级问题。
前置知识:
- Go 基础:goroutine、channel、
sync.Map/sync.RWMutex、context.Context、Redis 基本命令。 - Kratos 框架基础:
kratos.App启动流程、transport.Server接口。 robfig/cron/v3基本用法(知道 CRON 表达式格式即可)。
本章你会动手做的事:
- 用
cron.New(cron.WithSeconds(), cron.WithChain(cron.SkipIfStillRunning(...), cron.Recover(...)))起一个带防重叠 + panic 恢复的调度器。 - 给
wrapFunc加一层 RedisTryLock,模拟开两个进程,验证只有抢到锁的那个节点真正执行任务。 - 给存储校准任务加
ctx.Done()监听与"按user_id % N分片",体验优雅退出与并行加速。
一、技术栈与中间件
定时任务模块底层用 robfig/cron/v3 做调度引擎,并通过 Kratos 的 transport.Server 接入应用生命周期。下表汇总用到的技术与中间件,并标注了商用化需要补强的地方:
| 技术 / 中间件 | 用途说明 |
|---|---|
github.com/robfig/cron/v3 | Go 生态主流调度库,支持标准与秒级(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.status 与 StartAll/StopAll 临界区。 |
context.Context | 任务函数签名 func(ctx context.Context) error 的取消与超时传递;StopAll 时 cancel() 通知在跑任务退出。 |
transport.Server(Kratos) | 把调度器适配成"服务",复用框架优雅启停。 |
github.com/google/wire | 编译期依赖注入,组装 TaskManager 与 ScheduledTaskServer。 |
log/slog | 结构化日志(开始 / 完成 / 失败 / 耗时)。 |
data/lock 分布式锁(项目已有) | (商用必加) 基于 Redis SET NX + Lua 的 TryLock,带 TTL 自动释放与后台续期(renewLoop)。用于多实例下保证"同任务单点执行"。 |
| Prometheus 指标(需补) | (商用必加) 记录任务耗时、成功率、最近成功时间戳,供 Grafana / AlertManager 监控告警。 |
二、实现思路流程(总体)
定时任务模块遵循"分层抽象 + 生命周期托管 + 单点执行保护"的设计,从底层调度器到业务任务逐层封装:
cron 调度器初始化(SkipIfStillRunning + Recover + 时区)
NewCronTaskManager通过cron.New(cron.WithSeconds(), cron.WithChain(cron.SkipIfStillRunning(...), cron.Recover(...)))创建调度器。修正点:商用版需追加cron.WithLocation(time.LoadLocation("Asia/Shanghai"))固定时区。
任务定义(ScheduledTask 结构)
- 定义
ScheduledTask接口(Name/Spec/Func/Status/Start/Stop/Restart)与TaskFunc(func(ctx context.Context) error),NewTask便捷构造。
- 定义
任务注册(sync.Map 管理)
cronTaskManager用sync.Map维护任务名 -> *taskWrapper;Register去重(同名返回ErrTaskAlreadyExists),但不立即加入 cron,需等Start/StartAll`。
单点执行保护(分布式锁,商用必加)
- 修正点:在多实例部署下,每个 pod 都会各自跑一遍 cron。必须在任务被触发时先用 Redis
TryLock("cloud-disk:scheduler🔒{name}")抢占,抢不到就跳过;否则回收站清理会被多节点并发执行,导致存储被重复扣减、事件被重复发布。
- 修正点:在多实例部署下,每个 pod 都会各自跑一遍 cron。必须在任务被触发时先用 Redis
任务启动 / 优雅停机
StartAll通过cron.AddFunc注册并cron.Start();wrapFunc统一做"状态流转 + 耗时统计 + 锁/可观测/告警"。StopAll调cron.Stop()并在ctx上cancel(),最多等 10 秒让在跑任务退出。
ScheduledTaskServer 集成 Kratos 生命周期
ScheduledTaskServer实现transport.Server,Start→StartAll、Stop→StopAll,随app.Run()启停,无需业务代码手动管理。
具体业务任务
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 newTaskManager → Register 记账 → ScheduledTaskServer.Start(Kratos 启动)→ StartAll → cron.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.go 的 TryLock,带 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/Shanghai,0 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_at(cloud-disk:scheduler:last:{name}); - 进程启动
StartAll时,若now - last_success > 调度周期且错过了窗口,立即补跑一次; - 或对强一致任务用 K8s CronJob(自带
startingDeadlineSecondsbackfill)。
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.Locker,wrapFunc 触发时先 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)
NewCronTaskManager 传 loc *time.Location;StartAll 时读取 last_success,对"错过窗口"的任务立即补跑。Recover/锁异常时不影响补跑判定。
5. 任务分片并行 + 限流
存储校准从"单 goroutine 串行遍历全量用户"改为按 user_id % N 分片、N 个 worker 并行;对下游 DB / 对象存储调用用 semaphore.Weighted 或 rate.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)让所有节点"在同一时刻"触发,消除容器时区漂移。 - 修正点:
locker与metrics通过Option注入,既满足 wire 兼容,又把商用能力挂到wrapFunc。
5.2 任务定义与注册(ScheduledTask + TaskManager)
实现思路
scheduler.go 定义 ScheduledTask 接口、TaskManager 接口、TaskFunc 签名、TaskStatus 枚举、NewTask + simpleTask。taskWrapper 在 simpleTask 上增加 cron.EntryID 与带锁 status。Register 用 sync.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.go 中 taskWrapper 的 status 用读写锁保护(真正并发安全版),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/StartAll 才 AddFunc,保证"先批量注册、再统一启动"的状态一致。
5.3 ScheduledTaskServer 生命周期与优雅停机
实现思路
ScheduledTaskServer 实现 Kratos transport.Server:Start→StartAll、Stop→StopAll。StopAll 先 cron.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
}
要点:修正点:MarkExpiredDeleted 用 UPDATE ... 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)
这是把"商用审查重点"一次性落地的核心。wrapFunc 在 Recover 之内、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,但只在访问时判 ErrShareExpired(share.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
},
)
}
要点:过期分享不主动清理只会造成"分享表无限膨胀 + 用户看到已失效的分享",属于商用版应补的定时任务。其实现同样遵循"幂等 + 分布式锁"原则。
自测题与动手练习
自测题(合上书能答出来,才算懂):
SkipIfStillRunning和Recover各自解决什么?在WithChain里包裹顺序如何影响行为?它们能防"多实例重复执行"吗?- 多副本部署下,本模块原版的三个任务各节点都会跑。哪个任务因此会重复扣减存储?项目已有但未接入的什么能力可以解决它?请画出抢锁流程图。
- 回收站清理要做到"即使没锁也不重复扣存储",业务层该怎么改?(提示:乐观锁置位
status=deleted) robfig/cron在进程宕机期间错过的调度会补跑吗?商用如何做 backfill?请用一句话说明。NewCronTaskManager()原版用默认本地时区有什么风险?容器镜像默认时区常是什么?正确做法?wrapFunc原版只在失败记日志。生产应加哪三件事(重试 / 告警 / 死信),并区分哪两类错误?- 存储校准任务要支持"优雅退出 + 并行加速",循环里要监听什么?按什么字段分片?用什么限制下游并发?
- 为什么存储校准本身是幂等的,而回收站清理原版不是?分别说明原因。
动手练习(建议真做一遍):
- 起
cron.New(cron.WithSeconds(), cron.WithChain(cron.SkipIfStillRunning(cron.DiscardLogger), cron.Recover(cron.DefaultLogger))),写个 1 秒一次、故意time.Sleep(3s)的任务,观察是否被跳过。 - 给
wrapFunc加一层 RedisTryLock("cloud-disk:scheduler🔒demo"),开两个进程跑同一任务,验证只有抢到锁的节点真正执行、另一个记日志跳过。 - 把存储校准从"串行遍历"改成"按
user_id % 8分 8 个 worker +semaphore",循环里加if ctx.Err() != nil { break },跑起来后StopAll验证能优雅退出而非强杀。 - 给回收站清理加
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 指标(耗时 / 最近成功时间戳 / 失败计数),使"任务多久没跑、跑得快慢、是否连续失败"可被监控。 - 优雅停机与并行:
StopAll用cron.Stop()+cancel(ctx)+ 10s 等待;长任务循环监听ctx.Done(),存储校准按user_id % N分片并行 +semaphore限流。 - 缺失能力:过期分享仅访问时判过期,应补
share-clean定时任务(同样遵循幂等 + 分布式锁)。 - 下一篇可深入"分布式任务调度框架选型"(XXL-JOB / Elastic-Job / asynq / K8s CronJob),把本模块的抽象真正扩展到多机集群调度。