From f794cb43f0b64e14ef3371ed85954efb8aac3828 Mon Sep 17 00:00:00 2001 From: Jian Qiu Date: Fri, 7 Mar 2025 01:50:33 +0800 Subject: [PATCH] fix: task clean job in batch limit (#22238) Co-authored-by: Qiu Jian --- pkg/cloudcommon/consts/db.go | 10 ++++++++++ pkg/cloudcommon/db/taskman/tasks.go | 4 ++++ pkg/cloudcommon/options/changes.go | 3 +++ pkg/cloudcommon/options/options.go | 5 +++-- pkg/cloudevent/service/service.go | 2 +- pkg/cloudid/service/service.go | 2 +- pkg/compute/service/service.go | 2 +- pkg/devtool/models/devtoolcronjob.go | 2 +- pkg/image/service/service.go | 2 +- pkg/keystone/service/service.go | 2 +- pkg/monitor/service/service.go | 2 +- pkg/notify/service/service.go | 2 +- 12 files changed, 28 insertions(+), 10 deletions(-) diff --git a/pkg/cloudcommon/consts/db.go b/pkg/cloudcommon/consts/db.go index 1611d5576c..34ca3edae3 100644 --- a/pkg/cloudcommon/consts/db.go +++ b/pkg/cloudcommon/consts/db.go @@ -32,6 +32,8 @@ var ( taskArchiveThresholdHours int + taskArchiveBatchLimit int + enableChangeOwnerAutoRename = false enableDefaultPolicy = true @@ -90,6 +92,10 @@ func SetTaskArchiveThresholdHours(hours int) { taskArchiveThresholdHours = hours } +func SetTaskArchiveBatchLimit(limit int) { + taskArchiveBatchLimit = limit +} + func TaskWorkerCount() int { return taskWorkerCount } @@ -101,3 +107,7 @@ func LocalTaskWorkerCount() int { func TaskArchiveThresholdHours() int { return taskArchiveThresholdHours } + +func TaskArchiveBatchLimit() int { + return taskArchiveBatchLimit +} diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index ccc9bbc5fb..02b045a690 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -1424,6 +1424,10 @@ func (manager *STaskManager) migrateObjectInfo() error { func (manager *STaskManager) TaskCleanupJob(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { q := manager.Query().LT("created_at", time.Now().Add(-time.Duration(consts.TaskArchiveThresholdHours())*time.Hour)).Asc("created_at") + if consts.TaskArchiveBatchLimit() > 0 { + q = q.Limit(consts.TaskArchiveBatchLimit()) + } + rows, err := q.Rows() if err != nil { log.Errorf("query rows fail %s", err) diff --git a/pkg/cloudcommon/options/changes.go b/pkg/cloudcommon/options/changes.go index e4e81dd9b8..3a5b58abb2 100644 --- a/pkg/cloudcommon/options/changes.go +++ b/pkg/cloudcommon/options/changes.go @@ -71,6 +71,9 @@ func OnBaseOptionsChange(oOpts, nOpts interface{}) bool { if oldOpts.TaskArchiveThresholdHours != newOpts.TaskArchiveThresholdHours { consts.SetTaskArchiveThresholdHours(newOpts.TaskArchiveThresholdHours) } + if oldOpts.TaskArchiveBatchLimit != newOpts.TaskArchiveBatchLimit { + consts.SetTaskArchiveBatchLimit(newOpts.TaskArchiveBatchLimit) + } if oldOpts.EnableChangeOwnerAutoRename != newOpts.EnableChangeOwnerAutoRename { consts.SetChangeOwnerAutoRename(newOpts.EnableChangeOwnerAutoRename) } diff --git a/pkg/cloudcommon/options/options.go b/pkg/cloudcommon/options/options.go index 51366d5fbb..d75b62b79d 100644 --- a/pkg/cloudcommon/options/options.go +++ b/pkg/cloudcommon/options/options.go @@ -74,8 +74,9 @@ type BaseOptions struct { TaskWorkerCount int `default:"4" help:"Task manager worker thread count, default is 4"` LocalTaskWorkerCount int `default:"4" help:"Worker thread count that runs local tasks, default is 4"` - TaskArchiveThresholdHours int `default:"168" help:"The threshold in hours to migrate tasks to archives, default is 7days(168hours)"` - TaskArchiveIntervalHours int `default:"2" help:"The interval to migrate tasks to archives, default is 2 hours"` + TaskArchiveThresholdHours int `default:"168" help:"The threshold in hours to migrate tasks to archive, default is 7days(168hours)"` + TaskArchiveIntervalMinutes int `default:"60" help:"The interval in mibutes to migrate tasks to archive, default is 1 hour"` + TaskArchiveBatchLimit int `default:"10000" help:"The maximal count of tasks to archivie in a batch, default is 10000"` DefaultProcessTimeoutSeconds int `default:"60" help:"request process timeout, default is 60 seconds"` diff --git a/pkg/cloudevent/service/service.go b/pkg/cloudevent/service/service.go index 420ea41a2d..e0ff959606 100644 --- a/pkg/cloudevent/service/service.go +++ b/pkg/cloudevent/service/service.go @@ -67,7 +67,7 @@ func StartService() { cron.AddJobEveryFewHour("AutoPurgeSplitable", 4, 30, 0, db.AutoPurgeSplitable, false) - cron.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + cron.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) cron.Start() defer cron.Stop() diff --git a/pkg/cloudid/service/service.go b/pkg/cloudid/service/service.go index 8c8082e223..0584efb3fa 100644 --- a/pkg/cloudid/service/service.go +++ b/pkg/cloudid/service/service.go @@ -97,7 +97,7 @@ func StartService() { cron.AddJobEveryFewHour("AutoPurgeSplitable", 4, 30, 0, db.AutoPurgeSplitable, false) - cron.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + cron.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) cron.Start() defer cron.Stop() diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index c5b530e678..3d2127b90f 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -215,7 +215,7 @@ func StartServiceWithJobsAndApp(jobs func(cron *cronman.SCronJobManager), appCll cron.AddJobEveryFewDays( "CleanRecycleDiskFiles", 1, 3, 0, 0, models.StoragesCleanRecycleDiskfiles, false) - cron.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + cron.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) if jobs != nil { jobs(cron) diff --git a/pkg/devtool/models/devtoolcronjob.go b/pkg/devtool/models/devtoolcronjob.go index 4b6d0d5efd..d30055f45e 100644 --- a/pkg/devtool/models/devtoolcronjob.go +++ b/pkg/devtool/models/devtoolcronjob.go @@ -120,7 +120,7 @@ func InitializeCronjobs(ctx context.Context) error { DevToolCronManager = cronman.InitCronJobManager(true, 8, options.Options.TimeZone) - DevToolCronManager.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + DevToolCronManager.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) DevToolCronManager.Start() Session := auth.GetAdminSession(ctx, "") diff --git a/pkg/image/service/service.go b/pkg/image/service/service.go index 805be2a114..2b12178a77 100644 --- a/pkg/image/service/service.go +++ b/pkg/image/service/service.go @@ -164,7 +164,7 @@ func StartService() { cron.AddJobEveryFewHour("AutoPurgeSplitable", 4, 30, 0, db.AutoPurgeSplitable, false) - cron.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + cron.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) cron.Start() } diff --git a/pkg/keystone/service/service.go b/pkg/keystone/service/service.go index 56a0ff6f02..fc02df2a97 100644 --- a/pkg/keystone/service/service.go +++ b/pkg/keystone/service/service.go @@ -120,7 +120,7 @@ func StartService() { cron.AddJobEveryFewHour("RemoveObsoleteInvalidTokens", 6, 0, 0, models.RemoveObsoleteInvalidTokens, true) - cron.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + cron.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) cron.Start() defer cron.Stop() diff --git a/pkg/monitor/service/service.go b/pkg/monitor/service/service.go index 83c60bfba8..d7b3465d03 100644 --- a/pkg/monitor/service/service.go +++ b/pkg/monitor/service/service.go @@ -90,7 +90,7 @@ func StartService() { //cron.AddJobAtIntervalsWithStartRun("MonitorResourceSync", time.Duration(opts.MonitorResourceSyncIntervalSeconds)*time.Minute*60, models.MonitorResourceManager.SyncResources, true) cron.AddJobEveryFewHour("AutoPurgeSplitable", 4, 30, 0, db.AutoPurgeSplitable, false) - cron.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + cron.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) cron.Start() defer cron.Stop() diff --git a/pkg/notify/service/service.go b/pkg/notify/service/service.go index b8e316aadc..29c496380d 100644 --- a/pkg/notify/service/service.go +++ b/pkg/notify/service/service.go @@ -85,7 +85,7 @@ func StartService() { cron.AddJobEveryFewHour("AutoPurgeSplitable", 4, 30, 0, db.AutoPurgeSplitable, false) cron.AddJobEveryFewDays("InitReceiverProject", 7, 0, 0, 0, models.InitReceiverProject, true) - cron.AddJobAtIntervals("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalHours)*time.Hour, taskman.TaskManager.TaskCleanupJob) + cron.AddJobAtIntervalsWithStartRun("TaskCleanupJob", time.Duration(options.Options.TaskArchiveIntervalMinutes)*time.Minute, taskman.TaskManager.TaskCleanupJob, true) cron.Start() }