首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >从零实现一个高效定时任务调度器:Go container/heap 优先队列实战

从零实现一个高效定时任务调度器:Go container/heap 优先队列实战

作者头像
技术圈
发布2026-07-21 12:09:36
发布2026-07-21 12:09:36
600
举报

在订单 30 分钟未支付取消等延时场景下,直接处理延时任务常见有两类做法:使用 time.Sleep 会为每个等待任务产生一个挂起的 Goroutine,并发达到百万级时会导致内存暴涨与上下文切换频繁;而使用 time.AfterFunc 虽然在等待期间不占用独立协程,但在海量并发下会产生庞大的 Go 运行时 timer 节点,显著拉高调度与垃圾回收(GC)堆扫描的压力。

最优雅的解法是:在内存中用**最小堆(Min-Heap)**统一维护到期时间,配合单个后台协程事件循环统一调度。

为什么选择最小堆

定时任务调度的核心诉求是:时刻以最低代价拿到最早到期的任务。

  • O(1) 获取堆顶:全局最早过期的任务始终处于堆顶。
  • O(\log N) 节点调整:插入新任务或弹出堆顶节点,上浮与下沉调整极快。
  • 连续内存与 GC 友好:基于 Go 切片(Slice)实现,结构紧凑对 CPU 缓存友好。

构建最小堆 TaskHeap

Go 标准库 container/heap 要求类型实现 heap.Interface。首先定义任务结构体与切片堆:

代码语言:javascript
复制
   type Task struct {
    ID        string
    ExecuteAt time.Time
    TaskFunc  func()
    index     int
}

type TaskHeap []*Task

实现 sort.Interface(比较到期时间戳,确保最早过期的任务处于堆顶):

代码语言:javascript
复制
   func (h TaskHeap)Len()int{
    returnlen(h)
}

func(h TaskHeap)Less(i, j int)bool{
    return h[i].ExecuteAt.Before(h[j].ExecuteAt)
}

func(h TaskHeap)Swap(i, j int){
    h[i], h[j]= h[j], h[i]
    h[i].index = i
    h[j].index = j
}

实现 PushPop(注意在 Pop 时将 task.index 重置为 -1 并置空 nil,防止堆取消操作误删节点与内存泄漏):

代码语言:javascript
复制
   func (h *TaskHeap)Push(x any){
    task := x.(*Task)
    task.index =len(*h)
    *h =append(*h, task)
}

func(h *TaskHeap)Pop() any {
    old :=*h
    n :=len(old)
    task := old[n-1]
    old[n-1]=nil
    task.index =-1// 标记已出堆,防止误删
    *h = old[0: n-1]
    return task
}

设计单协程事件循环

调度器结构体 Scheduler 维护并发锁、最小堆与任务输入通道:

代码语言:javascript
复制
   type Scheduler struct {
    mu        sync.Mutex
    tasks     TaskHeap
    addTaskCh chan *Task
    stopCh    chan struct{}
}

在锁内批量取出到期任务切片,解锁后再异步启动协程执行回调,避免持锁回调放大锁竞争:

代码语言:javascript
复制
   func (s *Scheduler)popExpiredTasks(now time.Time)(time.Duration,bool){
    s.mu.Lock()
    var expired []*Task
    for s.tasks.Len()>0&&!s.tasks[0].ExecuteAt.After(now){
        task := heap.Pop(&s.tasks).(*Task)
        expired =append(expired, task)
    }

    var d time.Duration
    var hasTasks bool
    if s.tasks.Len()>0{
        d = s.tasks[0].ExecuteAt.Sub(now)
        if d <0{
            d =0
        }
        hasTasks =true
    }
    s.mu.Unlock()

    for_, task :=range expired {
        if task.TaskFunc !=nil{
            go task.TaskFunc()
        }
    }

    return d, hasTasks
}

事件循环的主控制逻辑(基于 select 多路复用监听唤醒与新任务推入):

代码语言:javascript
复制
   func (s *Scheduler)eventLoop(){
    var timer *time.Timer
    for{
        d, hasTasks := s.popExpiredTasks(time.Now())
        if timer !=nil{
            timer.Stop()
        }
        if hasTasks {
            timer = time.NewTimer(d)
        }else{
            timer = time.NewTimer(time.Hour)
        }

        select{
        case<-s.stopCh:
            if timer !=nil{
                timer.Stop()
            }
            return
        case<-timer.C:
        case task :=<-s.addTaskCh:
            s.mu.Lock()
            heap.Push(&s.tasks, task)
            s.mu.Unlock()
        }
    }
}

定时任务调度器使用示例

外部业务只需初始化调度器,即可异步提交延迟任务:

代码语言:javascript
复制
   s := NewScheduler()
s.Start()
defer s.Stop()

// 提交 50ms 延迟任务
s.AddTask("order_123", 50*time.Millisecond, func() {
    fmt.Println("订单超时未支付,自动取消")
})

由于采用了最小堆自动排序,即使并发提交大量任务,调度器通常也能以较低开销按到期顺序触发。

生产落地避坑要点

  • 防 Panic 保护:判空 task.TaskFunc != nil 避免空闭包崩掉后台协程。
  • 锁范围内最小化:先在锁内批量取出到期节点,解锁后再并发启动 go task.TaskFunc(),防止任务回调逻辑与调度器形成死锁或锁竞争。
  • 取消任务与竞态处理AddTask 可返回 *Task 句柄,取消时在锁保护下校验 task.index != -1,并调用 heap.Remove(&s.tasks, task.index) 移除节点并将索引设为 -1
  • 数据持久化:内存堆作为单机高效触发引擎,底座可结合 Redis ZSet 防止服务重启丢失任务。
本文参与 腾讯云自媒体同步曝光计划,分享自微信公众号。
原始发表:2026-07-20,如有侵权请联系 cloudcommunity@tencent.com 删除

本文分享自 技术圈子 微信公众号,前往查看

如有侵权,请联系 cloudcommunity@tencent.com 删除。

本文参与 腾讯云自媒体同步曝光计划  ,欢迎热爱写作的你一起参与!

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • 为什么选择最小堆
  • 构建最小堆 TaskHeap
  • 设计单协程事件循环
  • 定时任务调度器使用示例
  • 生产落地避坑要点
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档