增加延时Task

This commit is contained in:
2025-09-13 12:25:27 +08:00
parent 2593097989
commit 4035172a4b
4 changed files with 136 additions and 15 deletions

View File

@@ -23,6 +23,9 @@ type Task interface {
// IsDone 检查任务是否已完成
IsDone() bool
// GetDescription 获取任务说明
GetDescription() string
}
// taskItem 任务队列中的元素
@@ -65,7 +68,7 @@ func (q *Queue) AddTask(task Task) {
priority: task.GetPriority(),
}
heap.Push(q.queue, item)
q.logger.Infow("任务已添加到队列", "任务ID", task.GetID())
q.logger.Infow("任务已添加到队列", "任务ID", task.GetID(), "任务描述", task.GetDescription())
}
// GetNextTask 获取下一个要执行的任务(优先级最高的任务)
@@ -78,7 +81,7 @@ func (q *Queue) GetNextTask() Task {
}
item := heap.Pop(q.queue).(*taskItem)
q.logger.Infow("从队列中获取任务", "任务ID", item.task.GetID())
q.logger.Infow("从队列中获取任务", "任务ID", item.task.GetID(), "任务描述", item.task.GetDescription())
return item.task
}
@@ -185,7 +188,7 @@ func (e *Executor) Stop() {
// SubmitTask 提交任务到执行器
func (e *Executor) SubmitTask(task Task) {
e.queue.AddTask(task)
e.logger.Infow("任务已提交", "任务ID", task.GetID())
e.logger.Infow("任务已提交", "任务ID", task.GetID(), "任务描述", task.GetDescription())
}
// worker 工作协程
@@ -203,13 +206,13 @@ func (e *Executor) worker(id int) {
// 获取下一个任务
task := e.queue.GetNextTask()
if task != nil {
e.logger.Infow("工作协程正在执行任务", "工作协程ID", id, "任务ID", task.GetID())
e.logger.Infow("工作协程正在执行任务", "工作协程ID", id, "任务ID", task.GetID(), "任务描述", task.GetDescription())
// 执行任务
if err := task.Execute(); err != nil {
e.logger.Errorw("任务执行失败", "工作协程ID", id, "任务ID", task.GetID(), "错误", err)
e.logger.Errorw("任务执行失败", "工作协程ID", id, "任务ID", task.GetID(), "任务描述", task.GetDescription(), "错误", err)
} else {
e.logger.Infow("任务执行成功", "工作协程ID", id, "任务ID", task.GetID())
e.logger.Infow("任务执行成功", "工作协程ID", id, "任务ID", task.GetID(), "任务描述", task.GetDescription())
}
} else {
// 没有任务时短暂休眠