Kratos 云盘项目:定时任务、可观测性与限流

2025-01-15T10:30:00+08:00 | 32分钟阅读 | 更新于 2025-01-15T10:30:00+08:00

@

学习目标

读完本文你应该能够:

  1. 说清 robfig/cron/v3WithSecondsWithLocation(Asia/Shanghai) 解决了什么问题,能徒手写一个固定时区、支持秒级 cron 的调度器。
  2. 解释多副本部署下,为什么同一任务只能有一个节点真正执行,以及本项目的分布式锁 recycle:clean:lock + 看门狗续期是怎么兜底的。
  3. 讲透"启动补跑(backfill)““防并发重叠(SkipIfStillRunning)““重试退避(runWithRetry 1s/2s/4s)““可重试判定(isTransient)“四件套,并能在白板上画出状态机。
  4. lastSuccessFailureCountalertFuncdeadLetterFunc 这四个钩子把"任务多久没跑 / 失败了多少次 / 失败了通知谁 / 失败落哪里"讲成一套可观测方案。
  5. 拆解 TimingMiddlewareRateLimitMiddlewareCORSMiddlewarerecoveryMiddleware 的设计取舍,尤其是 IP 限流的 X-Real-IP 取值优先级与中间件顺序。

前置知识:Go 基础语法、context.Context、Kratos 的 middleware.Middleware 责任链模型、sync.Mutex/sync.Map、Redis 分布式锁(SET NX + Lua)的基本概念。

动手 3 件

  • 把项目 make run 跑起来,把 recycle-clean 的 cron 临时改成 */30 * * * * *,观察 3 点之外的日志与 api.timing 打点。
  • wrapFunc 增加一条 Prometheus Counter 指标,记录每个任务的成功/失败次数,跑一遍 grafanaprometheus 看板。
  • RateLimitMiddleware 的内存 map 换成基于 Redis 的集群版滑动窗口,让多实例共享限流计数。

一、为什么需要定时任务

Q1. 云盘为什么一定要用定时任务?不能让用户手动点"清理"吗?

答: 先打个比方。你家的冰箱会自己除霜,而不是等你哪天想起来"今天该除霜了"再去按开关。云盘里有很多"到点就必须发生、但用户根本不关心"的脏活:回收站里的文件 7 天过期要物理删除、用户已用存储配额要定期与真实文件对账校准、过期的分享链接要失效、上传崩溃留下的孤儿分片要清理。这些事有两个共同特征——时间驱动而不是事件驱动,且必须由系统兜底而不是指望用户自觉。

如果交给用户手动触发,会出现三类问题:

  • 资源泄漏:用户删了文件进回收站就跑路,7 天后不清理,磁盘空间永远还不回来。
  • 数据不一致:前端展示"已用 1GB”,但后端因为秒传、分片残留实际占用了 1.2GB,没有定时校准就永远对不上。
  • 运维不可控:清理是重 IO 操作,必须放在凌晨低峰(本项目是每天凌晨 3 点),用户手动点根本选不了时间。

本项目的四个定时任务都集中在 cmd/server/main.gonewTaskManager 里注册:

func newTaskManager(userUC *biz.UserUsecase, recycleUC *biz.RecycleUsecase,
    shareUC *biz.ShareUsecase, storage biz.Storage, l lock.Lock) scheduler.TaskManager {
    manager := scheduler.NewCronTaskManager(scheduler.WithLocker(l))
    _ = manager.Register(newStorageCalibrateTask(userUC))   // 0 0 3 * * * 存储校准
    _ = manager.Register(newRecycleCleanTask(recycleUC))    // 0 0 3 * * * 回收站清理
    _ = manager.Register(newChunkCleanupTask(storage))      // 0 0 4 * * * 孤儿分片清理
    _ = manager.Register(newShareCleanTask(shareUC))        // 0 0 2 * * * 过期分享清理
    return manager
}

⚠️ 注意 NewCronTaskManager(scheduler.WithLocker(l))——分布式锁是从最外层注入的,而不是写死在业务里。这样单机测试可以传 localLock,生产传 redisLock,业务代码零改动。这是 DDD 里"依赖接口、实现可替换"的典型应用。


二、robfig/cron/v3 的基本功

Q2. 为什么用 robfig/cron/v3WithSecondsAsia/Shanghai 到底在解决什么?

答: 先把 cron 表达式本身类比成"闹钟的表盘”。标准 Linux cron 是 5 位(分 时 日 月 周),不带秒。但云盘这种业务有时候需要"每 30 秒检查一次孤儿分片”,5 位精度就不够了。robfig/cron/v3WithSeconds() 让表达式变成 6 位:秒 分 时 日 月 周

比如回收站清理任务是 "0 0 3 * * *"——第 1 位是秒=0,第 2 位分=0,第 3 位时=3,后面 * 表示每天每月每周。也就是"每天 03:00:00 触发”。

WithLocation(Asia/Shanghai) 解决的是容器时区漂移这个经典坑。Docker 基础镜像默认时区往往是 UTC,如果你写 "0 0 3",在 UTC 容器里其实是北京时间 11 点才跑,跟"凌晨 3 点低峰"的初衷完全违背。本项目在 NewCronTaskManager 里把时区写死:

// 修正点:固定时区,避免各机器本地时区不一致导致执行时间漂移。
loc, _ := time.LoadLocation("Asia/Shanghai")

c := cron.New(
    cron.WithSeconds(),                       // 6 位表达式,支持秒级
    cron.WithLocation(loc),                   // 固定 Asia/Shanghai,消除容器时区漂移
    cron.WithChain(
        cron.SkipIfStillRunning(cron.DiscardLogger), // 上一轮没跑完就跳过这一轮
        cron.Recover(cron.DefaultLogger),             // 任务 panic 不拖垮整个 cron
    ),
)

答(工程细节): cron.WithChain 把两个"拦截器"串在每次执行外面——SkipIfStillRunningRecover。它们和 Kratos 的中间件是一个思想:在真正的任务函数外面再包一层防御。后面第 4、第 5 问会展开。

⚠️ time.LoadLocation("Asia/Shanghai") 依赖系统的时区数据库(tzdata)。 Alpine 镜像经常缺这个,会静默回退到 UTC。本项目用 _,但生产环境建议用 time/tzdata 这个空导入把时区表打进二进制,或者用 loc, err := ... 显式处理错误,否则"时区漂移"的坑会以最隐蔽的方式出现。

Q2(续):调度器接口与 Kratos 生命周期是怎么接起来的?

答: robfig/cron 本身只是一个"定时器”,它不会自己跟着 Kratos 应用一起启动和优雅退出。本项目在 internal/data/scheduler/scheduler.go 里抽象出两层接口,把"任务定义"和"任务管理"解耦,再套一个 transport.Server 接入 Kratos 生命周期。

先说任务定义接口 ScheduledTask——任何一个可被调度的任务都要实现这五个方法:

type ScheduledTask interface {
    Name() string    // 任务唯一名,如 "recycle-clean"
    Spec() string    // CRON 表达式,如 "0 0 3 * * *"
    Func() TaskFunc  // 真正执行的函数
    Status() TaskStatus
    Start() error
    Stop() error
    Restart() error
}

其中 TaskFunc 就是 func(ctx context.Context) error。注意接口里 Start/Stop/Restart 的注释写着"controlled by TaskManager”——也就是说单个任务自己不决定何时跑,调度权在管理器手里。本项目提供两种实现:

  • taskWrapper:注册进 cronTaskManager 后由 cron 真正驱动,带了 cronIDsync.RWMutex 保护的 status
  • simpleTaskNewTask(name, spec, fn) 返回的轻量实现,给"随手注册一个任务"用(比如存储校准),它的 Start/Stop 只是改个内存状态,真正的调度还是靠管理器。

再说管理器接口 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
}

答(工程细节): 这套接口的最大价值是可测试cronTaskManager 实现了 TaskManager,但单测时可以 mock 一个内存版 TaskManager,根本不依赖真正的 cron 和 Redis,直接验证"注册后 ListTasks 能拿到、失败计数正确"等逻辑。Register 里用 sync.Map 存任务并防重名(ErrTaskAlreadyExists),Start 里用 cron.AddFunc(tw.spec, m.wrapFunc(tw)) 把任务挂到 cron 上,返回 cron.EntryID 存进 tw.cronID——后续 Stop 就靠 m.cron.Remove(tw.cronID) 摘掉。

最妙的是 ScheduledTaskServer,它把管理器包成 Kratos 的 transport.Server

type ScheduledTaskServer struct {
    manager TaskManager
}

func (s *ScheduledTaskServer) Start(ctx context.Context) error {
    return s.manager.StartAll() // 应用启动 → 所有定时任务开始调度
}
func (s *ScheduledTaskServer) Stop(ctx context.Context) error {
    return s.manager.StopAll()  // 应用关闭 → 优雅停止所有任务
}

这样在 cmd/server/main.go 里只要把 NewScheduledTaskServer(manager) 交给 app.Run(),Kratos 会在进程启动时自动 StartAll(注册任务 + 补跑错过的),在收到 SIGTERM 时自动 StopAllStopAll 里还有一段优雅退出:

ctx := m.cron.Stop() // 停止接收新触发,但已在跑的任务会继续
// ... 把 tw.cronID 清零、started 置 false、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")
}

注意 m.cron.Stop() 返回的是一个 channel,只有所有正在跑的任务都返回后这个 channel 才关闭;同时 cancel() 会通知 runWithRetry 里的 ctx.Done() 立刻停止退避等待。两者配合,实现"不再触发新的、正在跑的尽快收尾、最多等 10 秒"的优雅退出。这正是把 cron 接入框架生命周期的标准姿势,面试时讲清楚 transport.Server 这一个适配层,能体现你对"框架整合"的理解深度。

Q2(续二):四个定时任务各自负责什么,为什么时间错开?

答: 回顾 newTaskManager 注册的四兄弟,它们的时间被刻意错开,避免凌晨同时开跑把数据库和磁盘 IO 打满:

任务名cron职责重 IO 程度
share-clean0 0 2 * * *把过期分享链接置为失效低(只改状态)
storage-calibrate0 0 3 * * *遍历所有用户,用真实文件重新校准已用存储中(全表扫描)
recycle-clean0 0 3 * * *物理删除 7 天过期的回收站文件并扣减存储高(删磁盘+改库)
chunk-cleanup0 0 4 * * *清理上传崩溃残留的孤儿分片目录中(扫目录删文件)

storage-calibraterecycle-clean 都是 3 点,但一个只读对账、一个写删除,量级可控;最重的 recycle-clean 之后 1 小时再跑 chunk-cleanup,把峰值错开。storage-calibrate 的实现是 ListAllUserIDs 拿到全部用户,再逐个 RecalibrateStorage——这种"全量遍历"任务最怕中途 panic,所以外层 Recover 链和 runWithRetry 对它尤其重要。

⚠️ storage-calibrate 遍历全量用户是 O(N) 的,如果用户量到千万级,凌晨一次性扫全表会锁表很久。生产上应当改成"按 user_id 分段游标"或丢给离线任务(如 Spark)做,不要在单进程 cron 里同步全扫。本项目是练手项目,用户量小,可以接受。


三、多实例单点执行

Q3. 云盘上线肯定要部署多个副本做高可用,那凌晨 3 点的清理任务会同时在所有节点跑吗?

答: 这是定时任务上生产最容易翻车的地方。想象小区里有 3 个物业管家,闹钟都设成早上 6 点浇花。如果没有协调,3 个人 6 点一起去浇同一片花——浪费水事小,关键是"重复扣减存储配额"这种事会直接算错账。所以我们必须保证:集群里同一时刻只有一个节点真正执行清理,其余节点到点了也只是"看一眼,发现别人在干,自己回去睡觉”。

本项目的做法是分布式锁。锁的 key 在 wrapFunc 里拼出来:

func (m *cronTaskManager) wrapFunc(tw *taskWrapper) cron.FuncJob {
    return func() {
        tw.setStatus(TaskStatusRunning)

        // 修正点:多实例单点执行——抢不到锁直接跳过,避免重复执行(如回收站清理重复扣存储)。
        lockKey := "cloud-disk:scheduler🔒" + tw.name
        if m.locker != nil {
            // 可重入 + 长 TTL,配合 redisLock 后台续期,任务跑久也不怕锁过期被别节点抢走。
            got, err := m.locker.TryLock(m.ctx, lockKey,
                lock.WithReentrant(true), lock.WithTTL(30*time.Minute))
            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) }()
        }
        // ... 真正执行任务
    }
}

答(工程细节): 这里有三个关键点要讲给面试官听:

  1. TryLock 而不是 LockTryLock 抢不到立刻返回 got=false,不阻塞。定时任务的特点是"错过就错过,下一轮再来",绝不能为等锁而卡住 goroutine。
  2. 长 TTL(30 分钟)+ 看门狗续期:如果锁 TTL 只有 10 秒,但清理任务要跑 15 分钟,锁过期了别的节点就会抢到,又变成并发执行。所以本项目用了 WithTTL(30*time.Minute),并且 Redis 锁内部有后台 renewLoop 每隔 ttlSec/3(即 10 分钟)用 Lua 续一次期,只要任务还在跑,锁就不会掉。
  3. 可重入:同一 goroutine 再次加同一把锁不会死锁(虽然 wrapFunc 里实际不会重入,但复用 redisLock 的通用能力,避免后续调用方踩坑)。

下面是多实例单点执行的时序图,注意锁是定义在调度器外层wrapFunc),而不是业务里:

flowchart TD
    A[节点A 触发 recycle-clean] --> B{抢分布式锁
cloud-disk:scheduler🔒recycle-clean} C[节点B 触发 recycle-clean] --> B D[节点C 触发 recycle-clean] --> B B -->|抢到锁| E[执行清理
物理删除加扣减存储] B -->|未抢到| F[记录日志 直接跳过] E --> G[后台看门狗每10分钟续期] E --> H[任务结束 释放锁] G --> H

⚠️ 锁一定要在任务最外层加,而不是在 CleanExpired 业务里才加。本项目其实做了双重保险:调度器 wrapFunccloud-disk:scheduler🔒<name> 保证"只有一个节点进任务",业务 CleanExpired 又用 recycle:clean:lock 保证"即使调度器漏了,业务层也兜得住"。两层锁 key 不同但目的互补,面试时可以强调这种"纵深防御"思路。

CleanExpired 里的全局锁是这样的:

func (uc *RecycleUsecase) CleanExpired(ctx context.Context) error {
    // 全局清理锁:多副本部署时只有一个实例执行,避免重复物理删除与重复扣减
    if uc.locker != nil {
        if err := uc.locker.Lock(ctx, "recycle:clean:lock"); err != nil {
            return err
        }
        defer func() { _ = uc.locker.Unlock(ctx, "recycle:clean:lock") }()
    }
    // ... 后面展开
}

答(工程细节·锁的看门狗): 光有 30 分钟 TTL 还不够——如果任务真的跑了 40 分钟,锁还是会过期。本项目 redisLock 在抢到锁后会起一个后台 renewLoop 持续续期,这是 Redis 分布式锁的标准"看门狗"模式:

// renewLoop 每隔 ttlSec/3 用 Lua 续一次期,直到 stopRenew 被关闭。
func (l *redisLock) renewLoop(ctx context.Context, key string, token string, ttlSec int64, stop chan struct{}) {
    renewInterval := time.Duration(ttlSec/3) * time.Second // 30 分钟 TTL → 每 10 分钟续期
    if renewInterval < 1*time.Second {
        renewInterval = 1 * time.Second
    }
    ticker := time.NewTicker(renewInterval)
    defer ticker.Stop()
    for {
        select {
        case <-ticker.C:
            _, err := l.rdb.Eval(ctx, renewLuaScript, []string{key}, token, ttlSec).Result()
            if err != nil {
                log.Warn("lock renewal failed", "key", key, "err", err)
            }
        case <-stop:
            return
        }
    }
}

续期用 Lua 脚本保证"只有持锁者(token 匹配)才能续",避免 A 的锁过期后 B 抢到,A 的看门狗又偷偷把 B 的锁续上的串味问题。Unlock 时关闭 stopRenew 通道,看门狗 goroutine 自然退出,不会泄漏。这套机制让我们敢把 TTL 设得比任务预期耗时更长,同时不会因为任务偶发变慢而丢锁。

⚠️ 看门狗依赖 ctx 存活。如果传入的 ctx 在任务结束前就被取消,renewLoop 里的 Eval 会因 ctx 取消而失败,续期停了,锁可能提前过期。所以传给 wrapFuncm.ctx 是整个管理器生命周期的 context(在 NewCronTaskManagercontext.WithCancel),而非每次任务的短时 ctx——这个区别很关键。


四、启动补跑与防并发重叠

Q4. 如果服务凌晨 2:55 崩溃,3:00 的清理任务没跑成,重启后这份"错过的任务"要不要补?

答: 要补,否则用户回收站里过期的文件就一直堆着。但补跑有个铁律:同一个任务,本进程曾经成功跑过的,不要再补。否则节点频繁重启就会反复补跑同一个任务,又回到"重复执行"的老问题。

本项目用 backfill 开关 + backfilled map 解决。启动时在 StartAll 里:

// 修正点:错过执行补跑——进程启动即补跑宕机期间错过且本进程从未成功过的任务。
// 由 wrapFunc 内的分布式锁保证多实例下只有一个节点真正执行。
if m.backfill {
    m.tasks.Range(func(key, value interface{}) bool {
        tw := value.(*taskWrapper)
        if _, ok := m.backfilled.LoadOrStore(tw.name, struct{}{}); ok {
            return true // 本进程已补跑过
        }
        if _, ok := m.lastSuccess.Load(tw.name); ok {
            return true // 本进程已成功执行过,不重复补跑
        }
        go func(t *taskWrapper) {
            m.logger.Info("scheduler: backfill run", "name", t.name)
            if err := t.fn(m.ctx); err != nil {
                m.logger.Error("scheduler: backfill failed", "name", t.name, "err", err)
            } else {
                m.lastSuccess.Store(t.name, time.Now())
            }
        }(tw)
        return true
    })
}

答(工程细节): 这段逻辑里有三个守卫:

  • backfilled:本进程生命周期内是否已经补跑过这个任务。防止 StartAll 被 Kratos 重启路径调用多次时重复补。
  • lastSuccess:本进程是否已经成功执行过这个任务。如果成功过,说明"错过的那次"其实已经被常规调度覆盖了,不必补。
  • 开 goroutine 异步补跑:func(t *taskWrapper) 传值而不是用闭包捕获循环变量,避免经典的循环变量捕获 bug;补跑失败只记日志,不阻塞主流程。

补跑流程如下:

flowchart TD
    A[进程启动 调用 StartAll] --> B{backfill 开关开启?}
    B -->|否| Z[仅启动正常调度]
    B -->|是| C{backfilled 已记录该任务?}
    C -->|是| Z
    C -->|否| D{lastSuccess 已存在?}
    D -->|是| E[本进程已成功过 跳过补跑]
    D -->|否| F[异步 goroutine 补跑该任务]
    F --> G[执行任务函数 fn]
    G --> H{成功?}
    H -->|是| I[写入 lastSuccess]
    H -->|否| J[记录错误日志 等待下一轮]
    I --> K[标记 backfilled]
    J --> K

Q5. 如果清理任务一跑就是 40 分钟,下一轮 3 点又触发了怎么办?

答: 这正是 SkipIfStillRunning 的用武之地。回到第 2 问 cron.WithChain(cron.SkipIfStillRunning(...))。它的语义是:上一次还没跑完,这一次到点直接跳过,而不是另起一个 goroutine 并发跑。对"清理物理文件 + 扣减存储"这种任务,并发跑意味着重复删除、重复扣减,绝对不能忍。

另一个链是 Recovercron.Recover(cron.DefaultLogger) 会在任务函数 panic 时 recover 住,打一条日志,然后只是这次任务失败,cron 调度循环继续活着。如果没有 Recover,一个任务的 panic 会沿着 goroutine 向上冒,直接把整个 cron 进程带崩,所有任务全停。

⚠️ SkipIfStillRunning 用的是"跳过"策略。它不是把任务排队,而是直接丢弃这一拍。对于清理类任务"晚一点跑没关系",跳过是对的;但如果你做的是"每分钟必须精确发一次对账",跳过就会丢数据,那种场景应该改成交互式队列或者拉长执行时间。面试时要能区分"可丢弃"和"必须执行"两类任务。


五、重试、退避与可重试判定

Q6. 任务执行失败了,是直接放弃还是重试?怎么避免"越重试越糟"?

答: 类比发微信:对方没网(瞬时故障)你重发几次他能收到;但消息内容本身违规被平台拦了(业务错误),你重发一百次也是被拦,纯属浪费。所以重试的前提是——错误得是可恢复的

本项目的 runWithRetry 最多重试 3 次,退避时间是 1s、2s、4s(指数退避,基数 2):

// runWithRetry 对瞬时错误做指数退避重试(最多 3 次),业务错误直接返回不重试。
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
}

答(工程细节): 这段代码有 4 个值得背下来的细节:

  1. 重试上限 maxRetry = 3:无限重试会把故障放大成雪崩。3 次是经验值,配合指数退避足够覆盖瞬时抖动又不会拖太久。
  2. 退避 1<<uint(i):第 0 次失败等 1s,第 1 次等 2s,第 2 次等 4s。1<<i 就是 2 的 i 次方,比乘法更地道,也避免了浮点。
  3. ctx 取消不重试if m.ctx.Err() != nil 在每次循环开头判断。进程收到终止信号时,StopAllcancel(),此时必须立刻停,不能还傻等退避。
  4. 退避期间也监听 ctx.Done()select 同时等 time.Afterctx.Done(),保证优雅退出时不卡在 sleep 上。

瞬时可重试 vs 业务不可重试isTransient 判定:

func isTransient(err error) bool {
    if err == nil {
        return false
    }
    if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
        return false // 进程关闭/超时,不重试
    }
    var netErr *net.OpError
    if errors.As(err, &netErr) {
        return true // 网络层错误,可重试
    }
    msg := strings.ToLower(err.Error())
    if strings.Contains(msg, "connection reset") ||
        strings.Contains(msg, "timeout") ||
        strings.Contains(msg, "broken pipe") ||
        strings.Contains(msg, "i/o timeout") {
        return true // 典型瞬时网络错误
    }
    return false // 其余(业务错误)不重试
}

答(工程细节): 这里用 errors.Is 精确匹配 context.Canceled/DeadlineExceeded,用 errors.As 提取 *net.OpError,再用错误消息字符串兜底。注意一个反直觉点:context.DeadlineExceeded 不算瞬时可重试——任务是被上游超时砍掉的,重试大概率还是超时,所以直接返回。这正是"区分错误性质"的价值。

flowchart TD
    A[执行任务 fn] --> B{返回 nil?}
    B -->|是| OK[记录 lastSuccess 返回 nil]
    B -->|否| C{isTransient 为 true?}
    C -->|业务错误| FE[直接返回错误 不重试]
    C -->|瞬时错误| D{i 小于 3?}
    D -->|否 已达上限| FE
    D -->|是| E[指数退避
1s 然后 2s 然后 4s] E --> F{ctx 已取消?} F -->|是| CANCEL[返回 ctx.Err 停止] F -->|否| A

六、可观测:让定时任务"看得见"

Q7. 定时任务藏在后台跑,怎么知道它是活着还是已经挂了几天?

答: 后台任务最怕"静默失败"——代码没崩,但任务其实两周没跑成功了,谁都不知道。本项目围绕 cronTaskManager 内置了四个可观测抓手:

  1. lastSuccess:每个任务最近一次成功时间。监控侧只要算 now - lastSuccess,超过阈值(比如 26 小时没成功,意味着当天 3 点的那次没跑成)就告警。
  2. failures:每个任务的累计失败次数(atomic.Int64,并发安全)。
  3. alertFunc:失败时的告警钩子,可注入钉钉/企业微信/邮件。
  4. deadLetterFunc:失败落"死信"的钩子,可注入写 task_dead_letter 表,供事后复盘。

这四个都在 wrapFunc 执行完毕后统一处理:

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

if err == nil {
    // 修正点:可观测——记录最近成功时间,供"任务多久没跑"告警。
    m.lastSuccess.Store(tw.name, time.Now())
    tw.setStatus(TaskStatusRunning)
    m.logger.Info("scheduler: task completed", "name", tw.name, "duration", duration)
    return
}

// 修正点:失败计数 + 状态 + 告警 + 死信。
if v, _ := m.failures.LoadOrStore(tw.name, new(int64)); v != nil {
    atomic.AddInt64(v.(*int64), 1)
}
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)  // 写死信

对外暴露的两个查询方法:

// LastSuccess 返回任务最近一次成功执行时间,供监控/告警使用。
func (m *cronTaskManager) LastSuccess(name string) (time.Time, bool) {
    v, ok := m.lastSuccess.Load(name)
    if !ok {
        return time.Time{}, false
    }
    return v.(time.Time), true
}

// FailureCount 返回任务累计失败次数,供监控/告警使用。
func (m *cronTaskManager) FailureCount(name string) int64 {
    v, ok := m.failures.Load(name)
    if !ok {
        return 0
    }
    return atomic.LoadInt64(v.(*int64))
}

答(工程细节): 注意 failuressync.Map*int64,第一次用 LoadOrStore 放一个指针,之后用 atomic.AddInt64 累加。这样多个任务并发失败也不会丢计数,且不需要每次都加锁。告警与死信用函数钩子注入(WithAlertFunc / WithDeadLetterFunc),默认是空实现——不注入就不做事,注入了就钉钉/写表。这是选项模式(Option)的典型用法:

// WithAlertFunc 注入失败告警回调(如 Webhook / 钉钉 / 邮件)。
func WithAlertFunc(fn func(name string, err error)) Option {
    return func(m *cronTaskManager) { m.alertFunc = fn }
}
// WithDeadLetterFunc 注入死信回调(如写 task_dead_letter 表)。
func WithDeadLetterFunc(fn func(name string, err error)) Option {
    return func(m *cronTaskManager) { m.deadLetterFunc = fn }
}

⚠️ 钩子本身(alertFunc/deadLetterFunc不能再用会 panic 的逻辑。它们是"失败处理"的兜底,如果告警通知自己再 panic,而 wrapFunc 这层没有 recover,就会把 cron 的调度 goroutine 带崩。生产上 alert/writeDeadLetter 内部最好自带 recover,或只在最外层 recoveryMiddleware 之外再包一层守护。


七、幂等:重复执行也不能出错

Q8. 即使有分布式锁,清理任务会不会还是重复删同一个文件、重复扣存储?

答: 会,因为分布式锁不是 100% 严密的(极端情况下看门狗续期失败、Redis 主从切换丢锁)。所以业务层必须自己做幂等——即便任务被重复调用,结果也和调用一次一样。本项目的 CleanExpired 用了两道防线:

  1. 全局锁 recycle:clean:lock(上一问已讲,纵深防御)。
  2. 每项删除前的二次校验:遍历过期列表时,对每一个 item 先加一把 recycle:item:<id> 的细粒度锁,再加锁后重新查一次"这东西还在不在",不在就当已处理。
expired, err := uc.recycleRepo.ListExpired(ctx, time.Now())
// ...
for _, item := range expired {
    // 每项的幂等保护:删除前再查一次,若已被其他实例清理则跳过
    key := fmt.Sprintf("recycle:item:%d", item.ID)
    lockErr := uc.withLock(ctx, key, func() error {
        if _, ferr := uc.recycleRepo.FindByID(ctx, item.ID); ferr != nil {
            // 已不存在,视为已处理
            return nil
        }
        ids := uc.permanentlyDeleteItem(ctx, item.UserID, item)
        if len(ids) > 0 {
            deletedByUser[item.UserID] = append(deletedByUser[item.UserID], ids...)
        }
        if derr := uc.recycleRepo.Delete(ctx, item.ID); derr != nil {
            log.Error("recycle: failed to delete expired recycle item", "id", item.ID, "err", derr)
        }
        return nil
    })
    // ...
}

答(工程细节): permanentlyDeleteItem 内部还有一个精妙的幂等设计——先删磁盘物理文件,成功后才改数据库状态并扣减存储

// 合规物理删除:先确保磁盘文件删除成功,再标记数据库与扣减存储,
// 否则跳过本项(保留 recycle 记录以便下次定时任务重试),避免计数错乱。
if err := uc.storage.Delete(ctx, file.Path); err != nil {
    log.Error("recycle: failed to delete physical file; skip db update to allow retry", ...)
    return deletedFileIDs
}

这个顺序很关键:如果反过来"先扣存储再删文件",一旦删文件失败,用户的配额被扣了但文件还在,下次重试又会再扣一次,存储计数就错乱了。先物理后逻辑,保证了"物理文件没了才算数",重试安全。

⚠️ 幂等不是"加了锁就万事大吉"。分布式锁防的是"并发",幂等防的是"重复执行"(包括补跑、手动重试、消息重投)。两者正交,必须同时做。面试时把"全局锁 + 细粒度锁 + 二次查存在 + 先物理后逻辑"四层串起来讲,会非常加分。


八、中间件:请求的"安检流水线"

Kratos 的中间件本质是函数套函数的责任链:middleware.Middleware 类型是 func(Handler) Handler,每一层都能在调用下一个 Handler 前后插入逻辑。本项目在 internal/server/middleware.go 实现了 Timing、Recovery、Validate、JWT、CORS、CacheControl、RateLimit 七种。

Q9. TimingMiddleware 是干什么的?为什么要把 /healthz 排除?

答: 它记录每个请求从进到出的耗时,打一条 api.timing 日志,是性能监控最朴素也最有效的一招。代码很短:

func TimingMiddleware() middleware.Middleware {
    return func(handler middleware.Handler) middleware.Handler {
        return func(ctx context.Context, req interface{}) (interface{}, error) {
            start := time.Now()
            op := "unknown"
            if tr, ok := transport.FromServerContext(ctx); ok {
                op = tr.Operation() // 取路由名,如 /api/v1/files
            }
            reply, err := handler(ctx, req)
            duration := time.Since(start)
            // 修正点:健康检查等高频端点跳过打点,避免日志噪音
            if op != "/healthz" && op != "/readyz" {
                log.Info("api.timing",
                    "op", op,
                    "duration_ms", duration.Milliseconds(),
                    "err", err,
                )
            }
            return reply, err
        }
    }
}

答(工程细节): 排除 /healthz/readyz 是因为 K8s 探针每一两秒就打一次,如果也打点,监控日志会被健康检查的噪声淹没,真正慢的业务接口反而看不清。transport.FromServerContext(ctx) 是 Kratos 从 context 里取当前请求元信息(操作名、请求头)的标准姿势——后面限流、CORS 都靠它。

Q10. RateLimitMiddleware 怎么做限流?为什么 IP 取 X-Real-IP 优先?

答: 它基于客户端 IP 做滑动窗口限流:每个 IP 维护一个窗口(起始时间 + 计数),窗口内请求数超过 maxRequests 就返回 429。最大的亮点是惰性清理——不用后台 goroutine 定时扫 map,而是在每次请求时顺带删掉过期的条目,从根本上避免了"后台 goroutine 泄漏/忘记停"这类问题。

func RateLimitMiddleware(maxRequests int, window time.Duration) middleware.Middleware {
    var (
        mu      sync.Mutex
        entries = make(map[string]*rateLimitEntry)
        lastCleanup = time.Now()
    )
    cleanupInterval := window / 2
    if cleanupInterval < time.Second {
        cleanupInterval = time.Second
    }
    return func(handler middleware.Handler) middleware.Handler {
        return func(ctx context.Context, req interface{}) (interface{}, error) {
            var clientIP string
            if tr, ok := transport.FromServerContext(ctx); ok {
                if ht, ok := tr.(interface{ Request() *http.Request }); ok {
                    // 修正点:优先取反向代理基于 $remote_addr 填充的 X-Real-IP,
                    // X-Forwarded-For 可被客户端伪造且含多级代理链,仅作兜底并取首跳。
                    clientIP = ht.Request().Header.Get("X-Real-IP")
                    if clientIP == "" {
                        if xff := ht.Request().Header.Get("X-Forwarded-For"); xff != "" {
                            if idx := strings.IndexByte(xff, ','); idx >= 0 {
                                clientIP = strings.TrimSpace(xff[:idx]) // 只取首跳,防伪造
                            } else {
                                clientIP = strings.TrimSpace(xff)
                            }
                        }
                    }
                    if clientIP == "" {
                        clientIP = ht.Request().RemoteAddr
                    }
                }
            }
            if clientIP == "" {
                clientIP = "unknown"
            }

            mu.Lock()
            now := time.Now()
            if now.Sub(lastCleanup) >= cleanupInterval {
                cutoff := now.Add(-window)
                for ip, entry := range entries {
                    if entry.windowStart.Before(cutoff) {
                        delete(entries, ip) // 惰性清理过期条目
                    }
                }
                lastCleanup = now
            }
            entry, exists := entries[clientIP]
            if !exists || now.Sub(entry.windowStart) >= window {
                entries[clientIP] = &rateLimitEntry{count: 1, windowStart: now}
                mu.Unlock()
                return handler(ctx, req)
            }
            entry.count++
            if entry.count > maxRequests {
                mu.Unlock()
                return nil, errors.TooManyRequests("RATE_LIMITED", "请求过于频繁,请稍后再试")
            }
            mu.Unlock()
            return handler(ctx, req)
        }
    }
}

答(工程细节): 三个面试高频点:

  1. IP 取值优先级 X-Real-IP > X-Forwarded-For 首跳 > RemoteAddrX-Real-IP 是反向代理(Nginx)用真实的 $remote_addr 填的,客户端改不了;X-Forwarded-For 是请求头,客户端可以伪造一串假 IP,所以本项目只取它的第一个(离服务端最近、由可信代理填的那个),并且仅作兜底。用 RemoteAddr 兜底是因为直连场景没有这两个头。
  2. 惰性清理cleanupInterval = window/2,只有距离上次清理超过半个窗口才扫一遍 map 删过期项。这样既防止 map 无限膨胀,又不需要常驻 goroutine,没有泄漏风险。
  3. sync.Mutex 保护 map:Go 的 map 不是并发安全的,mu.Lock() 包住所有读写。注意 return handler(...) 之前都先 mu.Unlock(),别把锁带进业务 Handler。

答(工程细节·算法选型): 本项目选的是"固定窗口 + 滑窗重置"的简化版滑动窗口,而不是令牌桶(token bucket)。原因很实在:令牌桶需要后台 goroutine 持续发令牌或每次请求时计算令牌数,实现更复杂;而本项目的场景是"单实例、按 IP 粗粒度挡刷量",固定窗口(每个 IP 一个窗口起点 + 计数)已经够用,且天然不需要后台线程。代价是窗口边界处可能有"双倍突发"(如窗口交替时两拍各放 maxRequests 个),但限流本就是"尽力而为"的防护,不是精确计费,这个代价可接受。如果以后要做"平滑限流"(如 Guava 的 RateLimiter 那样匀速放行),再换令牌桶或漏桶不迟。另外注意当 clientIP == "unknown"(完全取不到 IP)时本项目把所有未知请求当成同一个 key 限流——虽然会误伤,但比"不限"安全,生产上应当尽量保证反代一定填 X-Real-IP

滑动窗口判断流程:

flowchart TD
    A[收到请求] --> B[提取客户端 IP
X-Real-IP 优先] B --> C{距上次清理超过
半个窗口?} C -->|是| D[扫描并删除过期条目] C -->|否| E{存在该 IP 的 entry?} D --> E E -->|不存在 或 已超窗| F[新建窗口 count=1] E -->|存在且在窗口内| G[count 加一] G --> H{count 大于上限?} H -->|是| I[返回 429 限流] H -->|否| J[放行 handler] F --> J

Q11. recoveryMiddleware 为什么必须有?

答: Go 里一个 goroutine panic 没 recover,整个进程就挂。recoveryMiddleware 在最外层 recover,把 panic 转成 500 返回,保证单个请求的 bug 不会拖垮整个服务进程。Kratos 官方其实自带 recovery.Recovery(),本项目这里又写了一版,逻辑一致:

func recoveryMiddleware() middleware.Middleware {
    return func(handler middleware.Handler) middleware.Handler {
        return func(ctx context.Context, req interface{}) (reply interface{}, err error) {
            defer func() {
                if r := recover(); r != nil {
                    if e, ok := r.(error); ok {
                        err = e
                    } else {
                        err = errors.InternalServer("PANIC", "internal server error")
                    }
                    log.Error("panic recovered", "err", r)
                }
            }()
            reply, err = handler(ctx, req)
            return
        }
    }
}

答(工程细节): 注意 recover() 必须放在 defer 里,且 defer 函数要修改命名的返回值 err——这样才能把 panic 转成 error 返回。这是 Go 错误处理的固定套路,面试手写几乎必考。

Q12. CORSMiddleware 为什么要排在 JWTAuthMiddleware 之前?

答: 浏览器的 CORS 预检(OPTIONS 请求)不带 Authorization 头。如果 CORS 排在 JWT 后面,OPTIONS 预检会先撞上 JWT 中间件,因为没有 token 直接返回 401,浏览器就会认为跨域失败,导致前端所有接口都挂。所以 CORS 必须当"第一道门",先识别 OPTIONS 并短路返回,根本不进入后续认证链。

// CORSMiddleware 按白名单处理跨域,正确处理 OPTIONS 预检。
// allowedOrigins 为可信源列表(如 https://cloud.example.com),绝不接受 "*"。
// 必须排在 JWTAuthMiddleware 之前,否则浏览器 OPTIONS 预检不带 token 会被 JWT 直接 401,导致前端跨域全部失败。
func CORSMiddleware(allowedOrigins ...string) middleware.Middleware {
    allowSet := make(map[string]bool, len(allowedOrigins))
    for _, o := range allowedOrigins {
        allowSet[o] = true
    }
    return func(handler middleware.Handler) middleware.Handler {
        return func(ctx context.Context, req interface{}) (interface{}, error) {
            if tr, ok := transport.FromServerContext(ctx); ok {
                if ht, ok := tr.(interface{ Request() *http.Request }); ok {
                    r := ht.Request()
                    origin := r.Header.Get("Origin")
                    if origin == "" || allowSet[origin] {
                        if h, ok := tr.(khttp.Transporter); ok {
                            h.ReplyHeader().Set("Access-Control-Allow-Origin", origin)
                            // 仅在可信源开启凭据,绝不与 "*" 共存
                            if allowSet[origin] {
                                h.ReplyHeader().Set("Access-Control-Allow-Credentials", "true")
                            }
                            h.ReplyHeader().Set("Access-Control-Allow-Methods", "GET,POST,PUT,DELETE,OPTIONS")
                            h.ReplyHeader().Set("Access-Control-Allow-Headers", "Authorization,Content-Type")
                            h.ReplyHeader().Set("Access-Control-Max-Age", "600")
                        }
                        // 预检请求直接短路返回,不再进入 JWT 等后续中间件
                        if r.Method == http.MethodOptions {
                            return nil, nil
                        }
                    }
                }
            }
            return handler(ctx, req)
        }
    }
}

答(工程细节): 三个安全要点:

  • 白名单,绝不用 *allowSet 只放行配置的受信赖源。如果写 Access-Control-Allow-Origin: * 又想带 cookie(Allow-Credentials: true),浏览器会直接拒绝——两者互斥。本项目只对白名单源开启 credentials,安全且合规。
  • OPTIONS 短路:预检请求 return nil, nil,直接结束,不进 JWT。这正是"顺序"的价值。
  • Max-Age: 600:让浏览器缓存预检结果 10 分钟,减少重复预检开销。

Q13. ctxKey 为什么要用自定义类型,而不是直接用字符串 "user_id"

答: Go 的 context.WithValue 的 key 是 interface{},如果用裸字符串 "user_id" 当 key,任何包都能用同一个字符串往 context 里塞值,极容易发生 key 冲突——A 包写了 "user_id",B 包也用 "user_id" 但存的是另一种类型,取出来类型断言就炸。本项目用自定义类型 ctxKey 杜绝这个问题:

// ctxKey 是 context 键的自定义类型,避免使用裸字符串导致与其他包冲突。
type ctxKey string

const (
    ctxKeyUserID   ctxKey = "user_id"
    ctxKeyUsername ctxKey = "username"
    ctxKeyToken    ctxKey = "token"
)

// 注入
ctx = context.WithValue(ctx, ctxKeyUserID, claims.UserID)
// 取出
func CtxUserID(ctx context.Context) uint64 {
    if id, ok := ctx.Value(ctxKeyUserID).(uint64); ok {
        return id
    }
    return 0
}

答(工程细节): 自定义类型 ctxKey 即使底层字符串也是 "user_id",它的类型和别的包的 string 不同,所以 ctx.Value(ctxKeyUserID) 只会匹配用同类型同值写入的 key,天然隔离。对比有些项目(如 service.CtxUserID)直接用裸 "user_id" 字符串,一旦引入的第三方库也用这个字符串当 key,就会静默串味。面试时这一个小细节能体现"是否写过生产级 Go 代码"的功底。

Q14. 这些中间件应该怎么排顺序?

答: 顺序决定了"哪道门先拦"。本项目的合理顺序是:CORS → Recovery → RateLimit → Validate → JWT → Timing → Handler。CORS 最先处理预检;Recovery 包最外层兜底 panic;RateLimit 在认证前先挡掉恶意刷量(省得给未认证请求做 JWT 校验浪费 CPU);JWT 负责认证并把用户信息注入 context;Timing 最后包住业务逻辑统计耗时。责任链如下:

flowchart LR
    A[HTTP 请求] --> B[CORSMiddleware
处理预检与跨域] B --> C[RecoveryMiddleware
兜底 panic] C --> D[RateLimitMiddleware
IP 滑动窗口限流] D --> E[ValidateMiddleware
校验请求体] E --> F[JWTAuthMiddleware
鉴权并注入用户信息] F --> G[TimingMiddleware
统计耗时] G --> H[业务 Handler]

⚠️ 限流放 JWT 之前还是之后是个权衡。放之前能挡住未认证的刷量攻击(本项目选法);但如果你要做"按用户配额限流"而不是"按 IP 限流",就必须在 JWT 之后、拿到 userID 才能限。本项目限的是 IP,所以放 JWT 前更省资源。面试时讲清这个取舍,比背顺序更值钱。


九、面试延伸

Q15. robfig/cron 和 Java 的 Quartz 比,有什么优劣?真要做分布式调度用什么?

答: robfig/cron单机内存调度器——它只在当前进程里按时间表触发函数,本身不具备:集群协调、任务持久化、错过任务的统一补跑、可视化控制台、分片广播。所以它适合"任务逻辑简单、配合外部分布式锁就能跑"的场景(本项目就是)。Quartz 自带 JDBCJobStore,可以把任务落库、支持集群模式,但它是 Java 生态。

真正生产级的分布式调度平台(跨语言、带控制台、失败重试、分片、依赖编排)通常用:

  • xxl-job:轻量、有管理后台、GLUE 模式,国内接受度高。调度中心统一触发,执行器回调结果。
  • Elastic-Job(基于 ZooKeeper):支持分片,适合海量任务水平拆分到多机。

本项目的取舍是"用 cron 做触发 + 自研分布式锁做单点 + 自研 backfill/重试/可观测",好处是零额外组件依赖(只靠 Redis),坏处是没有统一控制台任务状态不持久化到专用库。如果团队要管几十上百个任务,建议直接上 xxl-job,把"触发"交给平台,“单点/幂等/重试"仍由业务保证。

Q16. 任务失败了到底怎么告警才不漏?

答: 本项目给了 alertFunc(即时通知)+ deadLetterFunc(落库复盘)双通道,这是标准做法:

  • 即时通道(钉钉/飞书/PagerDuty):用于"现在就有人要去看"的紧急失败。alertFunc 注入后,每次 wrapFunc 失败都调一次。注意要限流——同一任务连续失败别每分钟都轰运维,可以加"相同任务 N 分钟内只告警一次"的静默窗口。
  • 死信通道(写 task_dead_letter 表):用于"事后审计和手动重试”。即使即时通知被忽略,死信表永远留痕,FailureCount 也能从表里聚合。
  • 趋势告警(Prometheus + lastSuccess/FailureCount:监控侧定时拉 LastSuccess(name),如果 now - lastSuccess > 26h 说明当天没跑成,直接告警。这比"失败后通知"更前置——它能发现"任务根本没触发"这种更隐蔽的故障。

⚠️ 告警本身要可观测:如果 alertFunc 调钉钉但钉钉挂了,这个失败不能又把主流程带崩。alert/writeDeadLetter 内部必须自带 recover 或错误吞掉,它们是"旁路",绝不能影响任务本身的完成状态记录。


自测题与动手练习

自测题(口头回答):

  1. robfig/cron/v3 开了 WithSeconds() 后,cron 表达式是几位?"0 0 3 * * *" 分别表示什么?如果不设 WithLocation,容器默认时区可能带来什么后果?
  2. 多副本部署下,本项目用哪把锁保证"只有一个节点真正清理回收站"?TryLockLock 在这里为什么选 TryLock?锁的 TTL 设 30 分钟、看门狗每 10 分钟续期,是为了解决什么问题?
  3. 进程凌晨崩溃错过了 3 点的任务,重启后靠哪两个 map 防止"重复补跑"?backfill 的语义是什么?
  4. runWithRetry 最多重试几次?退避序列是多少?isTransient 为什么把 context.DeadlineExceeded 判为"不可重试"?
  5. CleanExpired 的幂等是怎么做的(至少说出三层)?为什么"先删物理文件、再改库扣存储"而不是反过来?

动手练习:

  1. recycle-clean 的 cron 临时改成 "*/20 * * * * *"(每 20 秒),起两个进程实例,观察日志里是不是只有一个节点打印 recycle clean task running,另一个打印 another node holds lock, skip
  2. wrapFunc 增加一条 Prometheus 指标:scheduler_task_total{name, result} 计数器,在成功/失败时 Inc(),然后用 promhttp 暴露 /metrics,验证 FailureCount 与指标一致。
  3. RateLimitMiddleware 的内存 map[string]*rateLimitEntry 改写成基于 Redis 的滑动窗口(用 ZSET 存时间戳、定期 ZREMRANGEBYSCORE 清理),让多实例共享同一份限流计数,并压测验证集群总 QPS 上限符合预期。

本章小结

本文从 Kratos 云盘项目出发,把"定时任务 + 可观测 + 限流中间件"串成了一条完整的生产线:

  • 触发层robfig/cron/v3WithSeconds 支持秒级、WithLocation(Asia/Shanghai) 消除时区漂移;SkipIfStillRunning + Recover 防重叠、防 panic 拖垮进程。
  • 单点层wrapFunccloud-disk:scheduler🔒<name> 分布式锁 + 30 分钟 TTL + 看门狗续期,保证多副本只有一个节点执行;CleanExpired 再用 recycle:clean:lock 做纵深防御。
  • 可靠性层backfill 在重启后补跑错过的任务(靠 backfilled/lastSuccess 防重复);runWithRetry 做最多 3 次、1s/2s/4s 指数退避重试;isTransient 区分瞬时错误可重试与业务错误直接返回;context 取消立即停手。
  • 可观测层lastSuccess/FailureCount 提供"多久没跑 / 失败几次"的量化信号,alertFunc/deadLetterFunc 钩子实现"即时通知 + 死信复盘"双通道。
  • 幂等层:全局锁 + 细粒度 recycle:item:<id> 锁 + 二次查存在 + “先物理后逻辑"删除顺序,保证重复执行结果一致。
  • 中间件层TimingMiddleware 统计耗时(排除 /healthz)、RateLimitMiddleware 基于 X-Real-IP 优先的 IP 滑动窗口 + 惰性清理、recoveryMiddleware 兜底 panic、CORSMiddleware 必须排在 JWT 前并白名单化、自定义 ctxKey 类型避免 context key 冲突。
复习提示:
  • 定时任务可靠性三板斧:分布式锁防多副本执行 + SkipIfStillRunning 防重叠 + Recover 防 panic 拖垮进程。
  • backfill 的精髓:利用 lastSuccess 时间戳定位漏跑的任务窗口,配合 backfilled 标记防止重复补跑。
  • 中间件链顺序决定成败:CORS → JWT → RateLimit → Recovery → Timing,任意顺序错乱都会导致预检失败或被绕过。
  • 指数退避重试200*(i+1)ms = 200ms → 400ms → 600ms,比固定间隔更能应对瞬时抖动。
面试官
“先物理后逻辑"的删除顺序是什么意思?为什么要这样设计?
候选人

两种删除策略
① 先逻辑后物理:先标记状态为"已删除",再异步清理物理文件
② 先物理后逻辑:先删除物理文件,再更新数据库状态

本项目采用"先物理后逻辑"的原因:

  • 如果先逻辑删除成功,但物理删除失败,会产生孤儿文件占用存储空间
  • 先物理删除确保资源真正释放,即使 DB 更新失败也可以重试(物理删除是幂等的)
  • 物理删除失败时,逻辑删除还没发生,可以通过补偿任务重新尝试

    关键约束:物理删除必须加细粒度锁 recycle:item:<id>,防止并发场景下重复删除或遗漏。面试中如果问到删除顺序,这个分析框架能让面试官看到你对一致性边界的理解。
复习提示:
  • 定时任务三保险:分布式锁防多副本重复执行 + SkipIfStillRunning 防重叠 + Recover 防 panic 拖垮进程,三层缺一不可。
  • backfill 补跑机制:重启后靠 lastSuccess 时间戳定位漏跑任务,配合 backfilled 标记防重复执行——这是调度系统的"可靠保证"核心。
  • 幂等设计:全局锁 + 细粒度 item 锁 + 二次查存在 + “先物理后逻辑"删除顺序,四重保障确保重复执行结果一致。
  • 中间件链顺序:CORS → JWT → RateLimit → Recovery → Timing,顺序错了会导致预检失败或限流绕过。
  • 下一篇讲数据库与 GORM——它是业务数据的持久化层,和定时任务的"触发→执行→持久化"链路紧密相关。

最后,cron 本质是单机调度器,任务多了、要可视化、要分片时应当迁移到 xxl-job / Elastic-Job 这类平台;但"单点执行、幂等、重试退避、可观测"这四件事,无论用什么框架都绕不开——它们才是定时任务能上生产的真正护城河。

About Me

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

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

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

目标

学AI,加油!加油!