mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-02 14:56:38 +08:00
优化任务管理
This commit is contained in:
@@ -111,3 +111,13 @@ func GetTaskMeta(taskType string) *TaskMeta {
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetTaskMetaByAsynqTask 根据 Asynq 任务名称获取元数据
|
||||
func GetTaskMetaByAsynqTask(asynqTask string) *TaskMeta {
|
||||
for _, t := range DispatchableTasks {
|
||||
if t.AsynqTask == asynqTask {
|
||||
return &t
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
+62
-36
@@ -218,25 +218,15 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// 从数据库加载执行记录
|
||||
execution, err := model.GetTaskExecutionByTaskID(ctx, taskID)
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "[TaskExecutor] 查询执行记录失败 taskID=%s: %v", taskID, err)
|
||||
// 执行记录不存在,仍然执行任务但不记录状态
|
||||
_, execErr := handler.Execute(ctx, t.Payload())
|
||||
if execErr != nil {
|
||||
span.SetStatus(codes.Error, execErr.Error())
|
||||
return execErr
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// 更新状态为 running
|
||||
// 加载或动态创建执行记录
|
||||
now := time.Now()
|
||||
execution.Status = model.TaskExecutionStatusRunning
|
||||
execution.StartedAt = &now
|
||||
if err := model.UpdateTaskExecution(ctx, execution); err != nil {
|
||||
logger.ErrorF(ctx, "[TaskExecutor] 更新执行状态失败 taskID=%s: %v", taskID, err)
|
||||
execution, err := getOrCreateTaskExecution(ctx, taskID, t, now)
|
||||
if err == nil && execution != nil && execution.TriggeredBy != "schedule" {
|
||||
execution.Status = model.TaskExecutionStatusRunning
|
||||
execution.StartedAt = &now
|
||||
if updateErr := model.UpdateTaskExecution(ctx, execution); updateErr != nil {
|
||||
logger.ErrorF(ctx, "[TaskExecutor] 更新执行状态失败 taskID=%s: %v", taskID, updateErr)
|
||||
}
|
||||
}
|
||||
|
||||
// 开始计时
|
||||
@@ -245,26 +235,69 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error {
|
||||
// 执行业务逻辑
|
||||
result, execErr := handler.Execute(ctx, t.Payload())
|
||||
|
||||
// 计算耗时
|
||||
// 计算耗时并归档记录
|
||||
duration := time.Since(start)
|
||||
finishTime := time.Now()
|
||||
|
||||
completeTaskExecution(ctx, execution, t, duration, finishTime, result, execErr, span)
|
||||
|
||||
if execution == nil && execErr != nil {
|
||||
span.SetStatus(codes.Error, execErr.Error())
|
||||
return execErr
|
||||
}
|
||||
|
||||
return execErr
|
||||
}
|
||||
|
||||
// getOrCreateTaskExecution 获取已有的任务执行记录,如果不存在则针对已知任务类型动态创建记录
|
||||
func getOrCreateTaskExecution(ctx context.Context, taskID string, t *asynq.Task, now time.Time) (*model.TaskExecution, error) {
|
||||
execution, err := model.GetTaskExecutionByTaskID(ctx, taskID)
|
||||
if err == nil {
|
||||
return execution, nil
|
||||
}
|
||||
|
||||
meta := GetTaskMetaByAsynqTask(t.Type())
|
||||
if meta == nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
execution = &model.TaskExecution{
|
||||
TaskID: taskID,
|
||||
TaskType: meta.AsynqTask,
|
||||
TaskName: meta.Name,
|
||||
Status: model.TaskExecutionStatusRunning,
|
||||
Retryable: meta.Retryable,
|
||||
MaxRetry: meta.MaxRetry,
|
||||
RetryCount: 0,
|
||||
Payload: string(t.Payload()),
|
||||
TriggeredBy: "schedule",
|
||||
StartedAt: &now,
|
||||
}
|
||||
|
||||
if createErr := model.CreateTaskExecution(ctx, execution); createErr != nil {
|
||||
logger.ErrorF(ctx, "[TaskExecutor] 动态创建执行记录失败 taskID=%s: %v", taskID, createErr)
|
||||
return nil, createErr
|
||||
}
|
||||
|
||||
return execution, nil
|
||||
}
|
||||
|
||||
// completeTaskExecution 完成并更新任务执行记录的状态和执行结果
|
||||
func completeTaskExecution(ctx context.Context, execution *model.TaskExecution, t *asynq.Task, duration time.Duration, finishTime time.Time, result *TaskResult, execErr error, span trace.Span) {
|
||||
if execution == nil {
|
||||
return
|
||||
}
|
||||
|
||||
execution.Duration = duration.Milliseconds()
|
||||
execution.FinishedAt = &finishTime
|
||||
|
||||
if execErr != nil {
|
||||
// 执行失败
|
||||
execution.Status = model.TaskExecutionStatusFailed
|
||||
execution.ErrorMessage = execErr.Error()
|
||||
|
||||
logger.ErrorF(ctx,
|
||||
"[TaskExecutor] 任务处理失败 Type: %s TaskID: %s Duration: %d ms Error: %v",
|
||||
t.Type(), taskID, duration.Milliseconds(), execErr,
|
||||
)
|
||||
|
||||
logger.ErrorF(ctx, "[TaskExecutor] 任务处理失败 Type: %s TaskID: %s Duration: %d ms Error: %v", t.Type(), execution.TaskID, duration.Milliseconds(), execErr)
|
||||
span.SetStatus(codes.Error, execErr.Error())
|
||||
span.RecordError(execErr)
|
||||
} else {
|
||||
// 执行成功
|
||||
execution.Status = model.TaskExecutionStatusSucceeded
|
||||
if result != nil {
|
||||
execution.Result = result.Message
|
||||
@@ -272,19 +305,12 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error {
|
||||
execution.Result = fmt.Sprintf("%s\n%s", result.Message, result.Detail)
|
||||
}
|
||||
}
|
||||
|
||||
logger.InfoF(ctx,
|
||||
"[TaskExecutor] 任务处理完成 Type: %s TaskID: %s Duration: %d ms",
|
||||
t.Type(), taskID, duration.Milliseconds(),
|
||||
)
|
||||
logger.InfoF(ctx, "[TaskExecutor] 任务处理完成 Type: %s TaskID: %s Duration: %d ms", t.Type(), execution.TaskID, duration.Milliseconds())
|
||||
}
|
||||
|
||||
// 更新执行记录
|
||||
if err := model.UpdateTaskExecution(ctx, execution); err != nil {
|
||||
logger.ErrorF(ctx, "[TaskExecutor] 更新执行记录失败 taskID=%s: %v", taskID, err)
|
||||
logger.ErrorF(ctx, "[TaskExecutor] 更新执行记录失败 taskID=%s: %v", execution.TaskID, err)
|
||||
}
|
||||
|
||||
return execErr
|
||||
}
|
||||
|
||||
// generateTaskID 生成任务 ID
|
||||
|
||||
@@ -1,67 +1,127 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package scheduler
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/logger"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
)
|
||||
|
||||
const (
|
||||
cleanupDedupWindow = 23 * time.Hour // 清理任务去重窗口
|
||||
cleanupMaxRetry = 3 // 清理任务最大重试次数
|
||||
)
|
||||
|
||||
var (
|
||||
scheduler *asynq.Scheduler
|
||||
schedulerOnce sync.Once
|
||||
activeScheduler *asynq.Scheduler
|
||||
schedulerMutex sync.Mutex
|
||||
quitChan chan struct{}
|
||||
schedulerOnce sync.Once
|
||||
)
|
||||
|
||||
func init() {
|
||||
// AsynqClient 已在 task 包中初始化
|
||||
}
|
||||
|
||||
// GetAsynqClient 获取全局 AsynqClient
|
||||
func GetAsynqClient() *asynq.Client {
|
||||
return task.AsynqClient
|
||||
}
|
||||
|
||||
// StartScheduler 启动调度器
|
||||
// StartScheduler 启动调度器 (该函数阻塞,直到调度器退出)
|
||||
func StartScheduler() error {
|
||||
var err error
|
||||
schedulerOnce.Do(func() {
|
||||
location, locErr := time.LoadLocation("Asia/Shanghai")
|
||||
if locErr != nil {
|
||||
err = fmt.Errorf(errLoadLocationFailed, locErr)
|
||||
return
|
||||
}
|
||||
scheduler = asynq.NewScheduler(
|
||||
task.RedisOpt,
|
||||
&asynq.SchedulerOpts{
|
||||
Location: location,
|
||||
},
|
||||
)
|
||||
quitChan = make(chan struct{})
|
||||
|
||||
// 清理未使用的上传文件任务
|
||||
if _, err = scheduler.Register(
|
||||
config.Config.Scheduler.CleanupUnusedUploadsTaskCron,
|
||||
asynq.NewTask(task.CleanupUnusedUploadsTask, nil),
|
||||
asynq.Unique(cleanupDedupWindow),
|
||||
asynq.MaxRetry(cleanupMaxRetry),
|
||||
); err != nil {
|
||||
// 初始化并运行首次调度
|
||||
if err = ReloadScheduler(); err != nil {
|
||||
err = fmt.Errorf("initial reload failed: %w", err)
|
||||
return
|
||||
}
|
||||
|
||||
// 启动调度器
|
||||
err = scheduler.Run()
|
||||
// 阻塞等待
|
||||
<-quitChan
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// StopScheduler 停止调度服务并解除 StartScheduler 阻塞
|
||||
func StopScheduler() {
|
||||
schedulerMutex.Lock()
|
||||
defer schedulerMutex.Unlock()
|
||||
|
||||
if activeScheduler != nil {
|
||||
activeScheduler.Shutdown()
|
||||
activeScheduler = nil
|
||||
}
|
||||
|
||||
if quitChan != nil {
|
||||
close(quitChan)
|
||||
quitChan = nil
|
||||
}
|
||||
}
|
||||
|
||||
// ReloadScheduler 重载调度器配置 (线程安全)
|
||||
func ReloadScheduler() error {
|
||||
schedulerMutex.Lock()
|
||||
defer schedulerMutex.Unlock()
|
||||
|
||||
// 1. 如果有运行中的调度器,先关闭它
|
||||
if activeScheduler != nil {
|
||||
activeScheduler.Shutdown()
|
||||
activeScheduler = nil
|
||||
}
|
||||
|
||||
// 2. 从数据库载入启用的定时任务配置
|
||||
schedules, err := model.ListActiveSchedules(context.Background())
|
||||
if err != nil {
|
||||
return fmt.Errorf("load schedules from db failed: %w", err)
|
||||
}
|
||||
|
||||
location, err := time.LoadLocation("Asia/Shanghai")
|
||||
if err != nil {
|
||||
return fmt.Errorf(errLoadLocationFailed, err)
|
||||
}
|
||||
|
||||
// 3. 实例化新的调度器
|
||||
newScheduler := asynq.NewScheduler(
|
||||
task.RedisOpt,
|
||||
&asynq.SchedulerOpts{
|
||||
Location: location,
|
||||
},
|
||||
)
|
||||
|
||||
// 4. 遍历并注册任务
|
||||
for _, s := range schedules {
|
||||
meta := task.GetTaskMeta(s.TaskType)
|
||||
if meta == nil {
|
||||
continue // 忽略排程配置中无效的任务类型
|
||||
}
|
||||
|
||||
// 构造 Asynq 载荷。定时任务使用对应 Meta 中的 Asynq 标识,同时将数据库中保存的 json 作为参数
|
||||
t := asynq.NewTask(meta.AsynqTask, []byte(s.Payload))
|
||||
|
||||
if _, err := newScheduler.Register(
|
||||
s.Cron,
|
||||
t,
|
||||
asynq.MaxRetry(meta.MaxRetry),
|
||||
asynq.Queue(meta.Queue),
|
||||
); err != nil {
|
||||
// 定时任务配置可能有误(如 Cron 格式不被 Asynq 识别),记录日志并跳过
|
||||
logger.ErrorF(context.Background(), "[Scheduler] 注册定时任务失败 id=%d name=%s: %v", s.ID, s.Name, err)
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
// 5. 替换全局调度器并异步启动
|
||||
activeScheduler = newScheduler
|
||||
go func() {
|
||||
if err := activeScheduler.Run(); err != nil {
|
||||
logger.ErrorF(context.Background(), "[Scheduler] 调度器运行错误: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
logger.InfoF(context.Background(), "[Scheduler] 成功重新加载定时任务,共注册 %d 个活动任务", len(schedules))
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user