MCPcopy Create free account
hub / github.com/cortexproject/cortex / cleanUser

Method cleanUser

pkg/compactor/blocks_cleaner.go:596–763  ·  view source on GitHub ↗
(ctx context.Context, userLogger log.Logger, userBucket objstore.InstrumentedBucket, userID string, firstRun bool)

Source from the content-addressed store, hash-verified

594}
595
596func (c *BlocksCleaner) cleanUser(ctx context.Context, userLogger log.Logger, userBucket objstore.InstrumentedBucket, userID string, firstRun bool) (returnErr error) {
597 c.blocksMarkedForDeletion.WithLabelValues(userID, reasonValueRetention)
598 startTime := time.Now()
599
600 bucketIndexDeleted := false
601
602 level.Info(userLogger).Log("msg", "started blocks cleanup and maintenance")
603 defer func() {
604 if returnErr != nil {
605 level.Warn(userLogger).Log("msg", "failed blocks cleanup and maintenance", "err", returnErr)
606 } else {
607 level.Info(userLogger).Log("msg", "completed blocks cleanup and maintenance", "duration", time.Since(startTime), "duration_ms", time.Since(startTime).Milliseconds())
608 }
609 c.tenantCleanDuration.WithLabelValues(userID).Set(time.Since(startTime).Seconds())
610 }()
611
612 if c.cfg.ShardingStrategy == util.ShardingStrategyShuffle && c.cfg.CompactionStrategy == util.CompactionStrategyPartitioning {
613 begin := time.Now()
614 c.cleanPartitionedGroupInfo(ctx, userBucket, userLogger, userID)
615 level.Info(userLogger).Log("msg", "finish cleaning partitioned group info files", "duration", time.Since(begin), "duration_ms", time.Since(begin).Milliseconds())
616 }
617
618 // Migrate block deletion marks to the global markers location. This operation is a best-effort.
619 if firstRun && c.cfg.BlockDeletionMarksMigrationEnabled {
620 if err := bucketindex.MigrateBlockDeletionMarksToGlobalLocation(ctx, c.bucketClient, userID, c.cfgProvider); err != nil {
621 level.Warn(userLogger).Log("msg", "failed to migrate block deletion marks to the global markers location", "err", err)
622 } else {
623 level.Info(userLogger).Log("msg", "migrated block deletion marks to the global markers location")
624 }
625 }
626
627 // Reading bucket index sync stats
628 idxs, err := bucketindex.ReadSyncStatus(ctx, c.bucketClient, userID, userLogger)
629
630 if err != nil {
631 level.Warn(userLogger).Log("msg", "error reading the bucket index status", "err", err)
632 idxs = bucketindex.Status{Version: bucketindex.SyncStatusFileVersion, NonQueryableReason: bucketindex.Unknown}
633 }
634
635 idxs.Status = bucketindex.Ok
636 idxs.SyncTime = time.Now().Unix()
637
638 // Read the bucket index.
639 begin := time.Now()
640 idx, err := bucketindex.ReadIndex(ctx, c.bucketClient, userID, c.cfgProvider, c.logger)
641
642 defer func() {
643 if bucketIndexDeleted {
644 level.Info(userLogger).Log("msg", "deleting bucket index sync status since bucket index is empty")
645 if err := bucketindex.DeleteIndexSyncStatus(ctx, c.bucketClient, userID); err != nil {
646 level.Warn(userLogger).Log("msg", "error deleting index sync status when index is empty", "err", err)
647 }
648 if err := c.deleteNonDataFiles(ctx, userLogger, userBucket); err != nil {
649 level.Warn(userLogger).Log("msg", "error deleting non-data files", "err", err)
650 }
651 } else {
652 bucketindex.WriteSyncStatus(ctx, c.bucketClient, userID, idxs, userLogger)
653 }

Calls 15

deleteNonDataFilesMethod · 0.95
EnableParquetMethod · 0.95
UpdateIndexMethod · 0.95
updateBucketMetricsMethod · 0.95
ReadSyncStatusFunction · 0.92
ReadIndexFunction · 0.92
DeleteIndexSyncStatusFunction · 0.92
WriteSyncStatusFunction · 0.92