package app import ( "context" "errors" "log/slog" "time" "hyapp/pkg/logx" ) func (a *App) startResourceEquipmentCleanupWorker(ctx context.Context) { a.workers.Add(1) go func() { defer a.workers.Done() ticker := time.NewTicker(a.resourceEquipmentCleanupWorkerCfg.PollInterval) defer ticker.Stop() for { if err := a.runResourceEquipmentCleanupRound(ctx); err != nil && !errors.Is(err, context.Canceled) { logx.Error(ctx, "wallet_resource_equipment_cleanup_round_failed", err) } select { case <-ctx.Done(): return case <-ticker.C: } } }() } func (a *App) runResourceEquipmentCleanupRound(ctx context.Context) error { if a.mysqlRepo == nil { return nil } cfg := a.resourceEquipmentCleanupWorkerCfg lockCtx, lockCancel := context.WithTimeout(ctx, cfg.QueryTimeout) release, acquired, err := a.mysqlRepo.AcquireResourceEquipmentCleanupLock(lockCtx) lockCancel() if err != nil || !acquired { return err } defer release() totalScanned := 0 totalDeleted := 0 totalDeadlockRetries := 0 reachedEnd := false for totalScanned < cfg.RoundMaxScannedRows { pageCtx, pageCancel := context.WithTimeout(ctx, cfg.QueryTimeout) result, pageErr := a.mysqlRepo.CleanupExpiredResourceEquipmentPage( pageCtx, cfg.ScanBatchSize, cfg.DeleteBatchSize, time.Now().UTC().UnixMilli(), ) pageCancel() if pageErr != nil { return pageErr } totalScanned += result.ScannedCount totalDeleted += result.DeletedCount totalDeadlockRetries += result.DeadlockRetries if result.ReachedEnd { reachedEnd = true break } if cfg.BatchPause > 0 { timer := time.NewTimer(cfg.BatchPause) select { case <-ctx.Done(): if !timer.Stop() { <-timer.C } return ctx.Err() case <-timer.C: } } } if totalDeleted > 0 || totalDeadlockRetries > 0 { // 日志只记录有实际回收或发生窄重试的轮次;空表巡检不制造周期性噪音。 logx.Info(ctx, "wallet_resource_equipment_cleanup_completed", slog.Int("scanned_count", totalScanned), slog.Int("deleted_count", totalDeleted), slog.Int("deadlock_retries", totalDeadlockRetries), slog.Bool("reached_end", reachedEnd), ) } return nil }