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

Function New

pkg/ingester/ingester.go:781–884  ·  view source on GitHub ↗

New returns a new Ingester that uses Cortex block storage instead of chunks storage.

(cfg Config, limits *validation.Overrides, registerer prometheus.Registerer, logger log.Logger, resourceMonitor *resource.Monitor)

Source from the content-addressed store, hash-verified

779
780// New returns a new Ingester that uses Cortex block storage instead of chunks storage.
781func New(cfg Config, limits *validation.Overrides, registerer prometheus.Registerer, logger log.Logger, resourceMonitor *resource.Monitor) (*Ingester, error) {
782 defaultInstanceLimits = &cfg.DefaultLimits
783 if cfg.ingesterClientFactory == nil {
784 cfg.ingesterClientFactory = client.MakeIngesterClient
785 }
786
787 bucketClient, err := bucket.NewClient(context.Background(), cfg.BlocksStorageConfig.Bucket, nil, "ingester", logger, registerer)
788 if err != nil {
789 return nil, errors.Wrap(err, "failed to create the bucket client")
790 }
791
792 i := &Ingester{
793 cfg: cfg,
794 limits: limits,
795 usersMetadata: map[string]*userMetricsMetadata{},
796 TSDBState: newTSDBState(bucketClient, registerer),
797 logger: logger,
798 ingestionRate: util_math.NewEWMARate(0.2, instanceIngestionRateTickInterval),
799 expandedPostingsCacheFactory: cortex_tsdb.NewExpandedPostingsCacheFactory(cfg.BlocksStorageConfig.TSDB.PostingsCache),
800 matchersCache: storecache.NoopMatchersCache,
801 }
802
803 if cfg.ActiveQueriedSeriesMetricsEnabled {
804 i.activeQueriedSeriesService = NewActiveQueriedSeriesService(logger, registerer)
805 }
806
807 if cfg.MatchersCacheMaxItems > 0 {
808 r := prometheus.NewRegistry()
809 registerer.MustRegister(cortex_tsdb.NewMatchCacheMetrics("cortex_ingester", r, logger))
810 i.matchersCache, err = storecache.NewMatchersCache(storecache.WithSize(cfg.MatchersCacheMaxItems), storecache.WithPromRegistry(r))
811 if err != nil {
812 return nil, err
813 }
814 }
815
816 i.metrics = newIngesterMetrics(registerer,
817 false,
818 cfg.ActiveSeriesMetricsEnabled,
819 cfg.ActiveQueriedSeriesMetricsEnabled,
820 i.getInstanceLimits,
821 i.ingestionRate,
822 &i.maxInflightPushRequests,
823 &i.maxInflightQueryRequests,
824 cfg.BlocksStorageConfig.TSDB.PostingsCache.Blocks.Enabled || cfg.BlocksStorageConfig.TSDB.PostingsCache.Head.Enabled,
825 cfg.EnableRegexMatcherLimits)
826 i.validateMetrics = validation.NewValidateMetrics(registerer)
827
828 // Replace specific metrics which we can't directly track but we need to read
829 // them from the underlying system (ie. TSDB).
830 if registerer != nil {
831 registerer.Unregister(i.metrics.memSeries)
832
833 promauto.With(registerer).NewGaugeFunc(prometheus.GaugeOpts{
834 Name: "cortex_ingester_memory_series",
835 Help: "The current number of series in memory.",
836 }, i.getMemorySeriesMetric)
837
838 promauto.With(registerer).NewGaugeFunc(prometheus.GaugeOpts{

Calls 14

NewClientFunction · 0.92
NewValidateMetricsFunction · 0.92
NewLifecyclerFunction · 0.92
NewFailureWatcherFunction · 0.92
NewResourceBasedLimiterFunction · 0.92
NewBasicServiceFunction · 0.92
newTSDBStateFunction · 0.85
newIngesterMetricsFunction · 0.85
NewLimiterFunction · 0.85
WatchServiceMethod · 0.80