
在订单 30 分钟未支付取消等延时场景下,直接处理延时任务常见有两类做法:使用 time.Sleep 会为每个等待任务产生一个挂起的 Goroutine,并发达到百万级时会导致内存暴涨与上下文切换频繁;而使用 time.AfterFunc 虽然在等待期间不占用独立协程,但在海量并发下会产生庞大的 Go 运行时 timer 节点,显著拉高调度与垃圾回收(GC)堆扫描的压力。
最优雅的解法是:在内存中用**最小堆(Min-Heap)**统一维护到期时间,配合单个后台协程事件循环统一调度。
定时任务调度的核心诉求是:时刻以最低代价拿到最早到期的任务。
Go 标准库 container/heap 要求类型实现 heap.Interface。首先定义任务结构体与切片堆:
type Task struct {
ID string
ExecuteAt time.Time
TaskFunc func()
index int
}
type TaskHeap []*Task
实现 sort.Interface(比较到期时间戳,确保最早过期的任务处于堆顶):
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
}
实现 Push 与 Pop(注意在 Pop 时将 task.index 重置为 -1 并置空 nil,防止堆取消操作误删节点与内存泄漏):
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 维护并发锁、最小堆与任务输入通道:
type Scheduler struct {
mu sync.Mutex
tasks TaskHeap
addTaskCh chan *Task
stopCh chan struct{}
}
在锁内批量取出到期任务切片,解锁后再异步启动协程执行回调,避免持锁回调放大锁竞争:
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 多路复用监听唤醒与新任务推入):
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()
}
}
}
外部业务只需初始化调度器,即可异步提交延迟任务:
s := NewScheduler()
s.Start()
defer s.Stop()
// 提交 50ms 延迟任务
s.AddTask("order_123", 50*time.Millisecond, func() {
fmt.Println("订单超时未支付,自动取消")
})
由于采用了最小堆自动排序,即使并发提交大量任务,调度器通常也能以较低开销按到期顺序触发。
task.TaskFunc != nil 避免空闭包崩掉后台协程。go task.TaskFunc(),防止任务回调逻辑与调度器形成死锁或锁竞争。AddTask 可返回 *Task 句柄,取消时在锁保护下校验 task.index != -1,并调用 heap.Remove(&s.tasks, task.index) 移除节点并将索引设为 -1。