fix: task clean job in batch limit (#22238)

Co-authored-by: Qiu Jian <qiujian@yunionyun.com>
This commit is contained in:
Jian Qiu
2025-03-07 01:50:33 +08:00
committed by GitHub
parent 9c0048dd05
commit f794cb43f0
12 changed files with 28 additions and 10 deletions
+10
View File
@@ -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
}
+4
View File
@@ -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)
+3
View File
@@ -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)
}
+3 -2
View File
@@ -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"`
+1 -1
View File
@@ -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()
+1 -1
View File
@@ -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()
+1 -1
View File
@@ -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)
+1 -1
View File
@@ -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, "")
+1 -1
View File
@@ -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()
}
+1 -1
View File
@@ -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()
+1 -1
View File
@@ -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()
+1 -1
View File
@@ -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()
}