diff --git a/CHANGELOG.md b/CHANGELOG.md index abd8bcc7927..1e021457f3c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ * [CHANGE] Querier: Make query time range configurations per-tenant: `query_ingesters_within`, `query_store_after`, and `shuffle_sharding_ingesters_lookback_period`. Uses `model.Duration` instead of `time.Duration` to support serialization but has minimum unit of 1ms (nanoseconds/microseconds not supported). #7160 * [CHANGE] Cache: Setting `-blocks-storage.bucket-store.metadata-cache.bucket-index-content-ttl` to 0 will disable the bucket-index cache. #7446 * [CHANGE] HA Tracker: Move `-distributor.ha-tracker.failover-timeout` from a global config to a per-tenant runtime config. The flag name and default value (30s) remain the same. #7481 +* [FEATURE] Ingester: Add experimental owned series tracking to prevent false throttling during ingester scale-up and ring resharding. The ingester counts only the series the ring assigns to it, rather than everything left in the TSDB head, closing a window of up to one head compaction cycle in which a tenant could be throttled against a limit that had already shrunk. Enabled via `-ingester.owned-series-metrics-enabled` and `-ingester.owned-series-limit-enforcement-enabled`. #7509 * [FEATURE] Parquet: Support sharded parquet file conversion and querying. #7610 * [FEATURE] Parquet Converter: Add experimental `-parquet-converter.max-num-columns` flag to automatically shard parquet files when the number of columns exceeds the configured limit. This prevents failures when a TSDB block has more unique label names than the parquet library's column limit (32767). #7624 * [FEATURE] Distributor: Add experimental `-distributor.num-query-workers` flag to use a goroutine worker pool for query fan-out calls to ingesters. Reuses pre-grown goroutine stacks to eliminate the `runtime.copystack` overhead (~8% CPU) observed on rulers with wide ingester fan-out. Falls back to spawning a new goroutine when no worker is available. #7623 diff --git a/docs/configuration/config-file-reference.md b/docs/configuration/config-file-reference.md index 0964ed55bed..689b3c7ee25 100644 --- a/docs/configuration/config-file-reference.md +++ b/docs/configuration/config-file-reference.md @@ -4198,6 +4198,18 @@ lifecycler: # CLI flag: -ingester.head-queried-series-metrics-windows [head_queried_series_metrics_windows: | default = 2h0m0s] +# Enable tracking of owned series per user. When enabled, the ingester computes +# series ownership based on the ring and emits cortex_ingester_owned_series +# metric. +# CLI flag: -ingester.owned-series-metrics-enabled +[owned_series_metrics_enabled: | default = false] + +# Use owned series count for limit enforcement. Requires +# owned-series-metrics-enabled. When enabled, PreCreation uses owned count +# instead of Head().NumSeries() for both per-user and instance-level limits. +# CLI flag: -ingester.owned-series-limit-enforcement-enabled +[owned_series_limit_enforcement_enabled: | default = false] + # Enable uploading compacted blocks. # CLI flag: -ingester.upload-compacted-blocks-enabled [upload_compacted_blocks_enabled: | default = true] diff --git a/docs/configuration/v1-guarantees.md b/docs/configuration/v1-guarantees.md index a957e03c8b1..92c82e1a51d 100644 --- a/docs/configuration/v1-guarantees.md +++ b/docs/configuration/v1-guarantees.md @@ -157,3 +157,7 @@ Currently experimental features are: - Parquet Converter: Maximum number of columns per file - `-parquet-converter.max-num-columns` (int) CLI flag - Automatically shards parquet files when the number of columns exceeds the configured limit +- Ingester: Owned Series Tracking + - Enable on Ingester via `-ingester.owned-series-metrics-enabled=true` + - Counts only the series the ring assigns to this ingester and exposes the `cortex_ingester_owned_series` metric + - `-ingester.owned-series-limit-enforcement-enabled` additionally uses that count for per-tenant and instance series limits, instead of the total number of series in the TSDB head diff --git a/pkg/distributor/distributor.go b/pkg/distributor/distributor.go index edf9b750b87..01bb9ee598b 100644 --- a/pkg/distributor/distributor.go +++ b/pkg/distributor/distributor.go @@ -36,7 +36,6 @@ import ( "github.com/cortexproject/cortex/pkg/ring" ring_client "github.com/cortexproject/cortex/pkg/ring/client" "github.com/cortexproject/cortex/pkg/util" - "github.com/cortexproject/cortex/pkg/util/extract" "github.com/cortexproject/cortex/pkg/util/flagext" "github.com/cortexproject/cortex/pkg/util/labelset" "github.com/cortexproject/cortex/pkg/util/limiter" @@ -587,49 +586,11 @@ func (d *Distributor) stopping(_ error) error { } func (d *Distributor) tokenForLabels(userID string, labels []cortexpb.LabelAdapter) (uint32, error) { - if d.cfg.ShardByAllLabels { - return shardByAllLabels(userID, labels), nil - } - - unsafeMetricName, err := extract.UnsafeMetricNameFromLabelAdapters(labels) - if err != nil { - return 0, err - } - return shardByMetricName(userID, unsafeMetricName), nil + return ring.TokenForLabels(userID, labels, d.cfg.ShardByAllLabels) } func (d *Distributor) tokenForMetadata(userID string, metricName string) uint32 { - if d.cfg.ShardByAllLabels { - return shardByMetricName(userID, metricName) - } - - return shardByUser(userID) -} - -// shardByMetricName returns the token for the given metric. The provided metricName -// is guaranteed to not be retained. -func shardByMetricName(userID string, metricName string) uint32 { - h := shardByUser(userID) - h = ingester_client.HashAdd32(h, metricName) - return h -} - -func shardByUser(userID string) uint32 { - h := ingester_client.HashNew32() - h = ingester_client.HashAdd32(h, userID) - return h -} - -// This function generates different values for different order of same labels. -func shardByAllLabels(userID string, labels []cortexpb.LabelAdapter) uint32 { - h := shardByUser(userID) - for _, label := range labels { - if len(label.Value) > 0 { - h = ingester_client.HashAdd32(h, label.Name) - h = ingester_client.HashAdd32(h, label.Value) - } - } - return h + return ring.TokenForMetadata(userID, metricName, d.cfg.ShardByAllLabels) } // Remove the label labelname from a slice of LabelPairs if it exists. diff --git a/pkg/distributor/distributor_test.go b/pkg/distributor/distributor_test.go index 2cb0b1522f3..b0317ee30b9 100644 --- a/pkg/distributor/distributor_test.go +++ b/pkg/distributor/distributor_test.go @@ -3799,7 +3799,7 @@ func (i *mockIngester) Push(ctx context.Context, req *cortexpb.WriteRequest, opt for j := range req.Timeseries { series := req.Timeseries[j] - hash := shardByAllLabels(orgid, series.Labels) + hash := ring.ShardByAllLabels(orgid, series.Labels) existing, ok := i.timeseries[hash] if !ok { // Make a copy because the request Timeseries are reused @@ -3818,7 +3818,7 @@ func (i *mockIngester) Push(ctx context.Context, req *cortexpb.WriteRequest, opt } for _, m := range req.Metadata { - hash := shardByMetricName(orgid, m.MetricFamilyName) + hash := ring.ShardByMetricName(orgid, m.MetricFamilyName) set, ok := i.metadata[hash] if !ok { set = map[cortexpb.MetricMetadata]struct{}{} @@ -4299,13 +4299,13 @@ func TestRemoveReplicaLabel(t *testing.T) { // This is not great, but we deal with unsorted labels when validating labels. func TestShardByAllLabelsReturnsWrongResultsForUnsortedLabels(t *testing.T) { t.Parallel() - val1 := shardByAllLabels("test", []cortexpb.LabelAdapter{ + val1 := ring.ShardByAllLabels("test", []cortexpb.LabelAdapter{ {Name: "__name__", Value: "foo"}, {Name: "bar", Value: "baz"}, {Name: "sample", Value: "1"}, }) - val2 := shardByAllLabels("test", []cortexpb.LabelAdapter{ + val2 := ring.ShardByAllLabels("test", []cortexpb.LabelAdapter{ {Name: "__name__", Value: "foo"}, {Name: "sample", Value: "1"}, {Name: "bar", Value: "baz"}, diff --git a/pkg/distributor/query.go b/pkg/distributor/query.go index 19e330e7533..c604897c747 100644 --- a/pkg/distributor/query.go +++ b/pkg/distributor/query.go @@ -104,7 +104,7 @@ func (d *Distributor) GetIngestersForQuery(ctx context.Context, matchers ...*lab metricNameMatcher, _, ok := extract.MetricNameMatcherFromMatchers(matchers) if ok && metricNameMatcher.Type == labels.MatchEqual { - return d.ingestersRing.Get(shardByMetricName(userID, metricNameMatcher.Value), ring.Read, nil, nil, nil) + return d.ingestersRing.Get(ring.ShardByMetricName(userID, metricNameMatcher.Value), ring.Read, nil, nil, nil) } } diff --git a/pkg/ingester/active_series.go b/pkg/ingester/active_series.go index 1a685a36cfa..e9945dc22f9 100644 --- a/pkg/ingester/active_series.go +++ b/pkg/ingester/active_series.go @@ -3,18 +3,70 @@ package ingester import ( "math" "sync" + "sync/atomic" "time" "github.com/prometheus/prometheus/model/labels" - "go.uber.org/atomic" + uatomic "go.uber.org/atomic" + + "github.com/cortexproject/cortex/pkg/ring" ) const ( numActiveSeriesStripes = 512 ) +// ringState holds the ring ownership data needed for series ownership checks. +// Stored behind an atomic.Pointer so that readers (hot push path) always see +// a consistent snapshot without lock contention, while the writer (periodic +// updateActiveSeries loop) can swap in a new state atomically. +type ringState struct { + // tokens is the ring's full sorted token list. + tokens []uint32 + // ownedPositions is parallel to tokens and marks the positions whose replica + // set includes this ingester. Both come from Lifecycler.GetOwnedTokenPositions + // and are immutable. + ownedPositions []bool +} + +// emptyRingState is the zero-value ring state used before any ring data is loaded. +var emptyRingState = &ringState{} + // ActiveSeries is keeping track of recently active series for a single tenant. +// +// It maintains two independent counts over the same set of entries: +// +// - active: entries whose last sample is newer than the idle timeout. This is +// the long-standing cortex_ingester_active_series gauge and its value is +// unchanged by owned-series tracking. +// - owned: entries whose ring token places them on this ingester, regardless of +// how recently they received a sample. +// +// The two use different retention, and that is deliberate. Entries live until +// Purge is called with the TSDB head's minimum time, so owned counts everything +// still pinned in the head. That makes owned a proxy for the memory the tenant +// is actually holding here, which is what a series limit is protecting. Because +// owned ignores the idle window while active applies it, owned may exceed active +// for a tenant with high churn. type ActiveSeries struct { + // Ring ownership state. Readers on the push path load atomically; + // the writer (updateTokens) stores a new pointer on ring changes. + ring atomic.Pointer[ringState] + + // currFingerprint detects ring changes. It is the fingerprint reported by the + // lifecycler alongside the ownership bitmap, which changes if and only if + // ownership may have changed. Only accessed by the updateTokens caller + // (periodic updateActiveSeries goroutine), so no synchronization needed. + currFingerprint uint64 + + // Cached totals across all stripes. PreCreation consults these on every new + // series, so they must not require locking 512 stripes to read. They are + // incremented as series are created and then recomputed authoritatively by + // UpdateMetrics and Purge, which also corrects any drift. + activeTotal uatomic.Int64 + ownedTotal uatomic.Int64 + activeNativeHistogramTotal uatomic.Int64 + stripes [numActiveSeriesStripes]activeSeriesStripe } @@ -23,23 +75,30 @@ type activeSeriesStripe struct { // Unix nanoseconds. Only used by purge. Zero = unknown. // Updated in purge and when old timestamp is used when updating series (in this case, oldestEntryTs is updated // without holding the lock -- hence the atomic). - oldestEntryTs atomic.Int64 + oldestEntryTs uatomic.Int64 mu sync.RWMutex refs map[uint64][]activeSeriesEntry - active int // Number of active entries in this stripe. Only decreased during purge or clear. - activeNativeHistogram int // Number of active entries only for Native Histogram in this stripe. Only decreased during purge or clear. + active int // Number of entries in this stripe within the idle window. + activeNativeHistogram int // Number of Native Histogram entries in this stripe within the idle window. + owned int // Number of entries in this stripe owned by this instance, ignoring the idle window. } // activeSeriesEntry holds a timestamp for single series. type activeSeriesEntry struct { - lbs labels.Labels - nanos *atomic.Int64 // Unix timestamp in nanoseconds. Needs to be a pointer because we don't store pointers to entries in the stripe. + lbs labels.Labels + key uint32 // Ring token hash for this series (used for ownership checks) + // owned caches whether key belongs to this instance, so that counting owned + // series does not repeat the ring lookup on every pass. Refreshed only when + // the ring changes. Guarded by the stripe mutex. + owned bool + nanos *uatomic.Int64 // Unix timestamp in nanoseconds. Needs to be a pointer because we don't store pointers to entries in the stripe. isNativeHistogram bool } func NewActiveSeries() *ActiveSeries { c := &ActiveSeries{} + c.ring.Store(emptyRingState) // Stripes are pre-allocated so that we only read on them and no lock is required. for i := range numActiveSeriesStripes { @@ -49,51 +108,149 @@ func NewActiveSeries() *ActiveSeries { return c } -// Updates series timestamp to 'now'. Function is called to make a copy of labels if entry doesn't exist yet. -func (c *ActiveSeries) UpdateSeries(series labels.Labels, hash uint64, now time.Time, nativeHistogram bool, labelsCopy func(labels.Labels) labels.Labels) { +// UpdateSeries updates series timestamp to 'now'. The key parameter is the ring token +// for this series (computed via ring.TokenForLabels). When the ring is not loaded, +// ownership is unknown and the series counts as owned. +// +// Every series is tracked. Ownership only decides whether the series counts +// towards owned, never whether it is tracked at all, so the active count behaves +// exactly as it did before owned-series tracking existed. +func (c *ActiveSeries) UpdateSeries(series labels.Labels, hash uint64, key uint32, now time.Time, nativeHistogram bool, labelsCopy func(labels.Labels) labels.Labels) { stripeID := hash % numActiveSeriesStripes - c.stripes[stripeID].updateSeriesTimestamp(now, series, hash, nativeHistogram, labelsCopy) + // Load ring state atomically — readers on the push path always see a consistent snapshot. + state := c.ring.Load() + created, owned := c.stripes[stripeID].updateSeriesTimestamp(now, series, hash, key, nativeHistogram, labelsCopy, state.tokens, state.ownedPositions) + + if !created { + return + } + + c.activeTotal.Inc() + if owned { + c.ownedTotal.Inc() + } + if nativeHistogram { + c.activeNativeHistogramTotal.Inc() + } +} + +// updateTokens updates the cached ring state. Returns true if ownership changed. +// Only called from the updateActiveSeries goroutine (single writer). +// +// The lifecycler's fingerprint is used rather than a hash of the token list, +// because an instance changing state alters the replica set, and therefore +// ownership, without changing any token. +func (c *ActiveSeries) updateTokens(tokens []uint32, ownedPositions []bool, fingerprint uint64) bool { + if len(tokens) == 0 || fingerprint == c.currFingerprint { + return false + } + + // The slices come from the lifecycler and are never mutated in place, so they + // can be published to readers directly rather than copied. + c.ring.Store(&ringState{ + tokens: tokens, + ownedPositions: ownedPositions, + }) + c.currFingerprint = fingerprint + + return true +} + +// UpdateMetrics recomputes the active and owned counts, re-evaluating ownership if +// the ring changed. Called from updateActiveSeries when OwnedMetrics is enabled. +// +// This deliberately removes nothing. Entries are released by Purge at head +// compaction instead, so that owned keeps counting series which are still in the +// head but have gone idle. keepUntil therefore only decides which entries count +// as active. +func (c *ActiveSeries) UpdateMetrics(keepUntil time.Time, tokens []uint32, ownedPositions []bool, fingerprint uint64) { + tokensChanged := c.updateTokens(tokens, ownedPositions, fingerprint) + + // Load the ring state from the atomic pointer for consistency. + // Even though we're on the same goroutine that just stored it, reading from + // the pointer ensures all code paths use the same access pattern. + state := c.ring.Load() + + var active, owned, activeNativeHistogram int64 + for s := range numActiveSeriesStripes { + a, o, nh := c.stripes[s].updateMetrics(keepUntil, tokensChanged, state.tokens, state.ownedPositions) + active += int64(a) + owned += int64(o) + activeNativeHistogram += int64(nh) + } + + c.activeTotal.Store(active) + c.ownedTotal.Store(owned) + c.activeNativeHistogramTotal.Store(activeNativeHistogram) } -// Purge removes expired entries from the cache. This function should be called -// periodically to avoid memory leaks. +// Purge removes entries last updated before keepUntil and recomputes the counts. +// +// When owned-series tracking is enabled this is called with the TSDB head's +// minimum time, so that an entry survives exactly as long as the series it +// describes is in the head. When tracking is disabled it is called with the idle +// timeout, which is the original behaviour. func (c *ActiveSeries) Purge(keepUntil time.Time) { + var active, owned, activeNativeHistogram int64 for s := range numActiveSeriesStripes { - c.stripes[s].purge(keepUntil) + a, o, nh := c.stripes[s].purge(keepUntil) + active += int64(a) + owned += int64(o) + activeNativeHistogram += int64(nh) } + + c.activeTotal.Store(active) + c.ownedTotal.Store(owned) + c.activeNativeHistogramTotal.Store(activeNativeHistogram) } -// nolint // Linter reports that this method is unused, but it is. +// clear drops every tracked entry. Used when the TSDB head is empty, so that no +// entry outlives the series it describes. func (c *ActiveSeries) clear() { for s := range numActiveSeriesStripes { c.stripes[s].clear() } + + c.activeTotal.Store(0) + c.ownedTotal.Store(0) + c.activeNativeHistogramTotal.Store(0) } +// Active returns the number of series which received a sample more recently than +// the idle timeout. Its value is not affected by ownership tracking. func (c *ActiveSeries) Active() int { - total := 0 - for s := range numActiveSeriesStripes { - total += c.stripes[s].getActive() - } - return total + return int(c.activeTotal.Load()) +} + +// Owned returns the number of tracked series whose ring token places them on this +// instance, ignoring the idle window. +// +// This may exceed Active for a tenant with high churn, because it also counts +// idle series which are still held in the TSDB head. That is intended: it is the +// count of series this instance is actually storing, which is what the series +// limit exists to bound. +// +// Before the ring has been read, ownership is unknown and every series counts as +// owned, so this equals Active. +func (c *ActiveSeries) Owned() int { + return int(c.ownedTotal.Load()) } func (c *ActiveSeries) ActiveNativeHistogram() int { - total := 0 - for s := range numActiveSeriesStripes { - total += c.stripes[s].getActiveNativeHistogram() - } - return total + return int(c.activeNativeHistogramTotal.Load()) } -func (s *activeSeriesStripe) updateSeriesTimestamp(now time.Time, series labels.Labels, fingerprint uint64, nativeHistogram bool, labelsCopy func(labels.Labels) labels.Labels) { +// updateSeriesTimestamp records a sample for a series, creating the entry if this +// is the first time it has been seen. It reports whether an entry was created and, +// if so, whether that entry is owned by this instance. +func (s *activeSeriesStripe) updateSeriesTimestamp(now time.Time, series labels.Labels, fingerprint uint64, key uint32, nativeHistogram bool, labelsCopy func(labels.Labels) labels.Labels, tokens []uint32, ownedPositions []bool) (created, owned bool) { nowNanos := now.UnixNano() e := s.findEntryForSeries(fingerprint, series) entryTimeSet := false if e == nil { - e, entryTimeSet = s.findOrCreateEntryForSeries(fingerprint, series, nowNanos, nativeHistogram, labelsCopy) + e, entryTimeSet, created, owned = s.findOrCreateEntryForSeries(fingerprint, key, series, nowNanos, nativeHistogram, labelsCopy, tokens, ownedPositions) } if !entryTimeSet { @@ -111,9 +268,11 @@ func (s *activeSeriesStripe) updateSeriesTimestamp(now time.Time, series labels. } } } + + return created, owned } -func (s *activeSeriesStripe) findEntryForSeries(fingerprint uint64, series labels.Labels) *atomic.Int64 { +func (s *activeSeriesStripe) findEntryForSeries(fingerprint uint64, series labels.Labels) *uatomic.Int64 { s.mu.RLock() defer s.mu.RUnlock() @@ -127,59 +286,102 @@ func (s *activeSeriesStripe) findEntryForSeries(fingerprint uint64, series label return nil } -func (s *activeSeriesStripe) findOrCreateEntryForSeries(fingerprint uint64, series labels.Labels, nowNanos int64, nativeHistogram bool, labelsCopy func(labels.Labels) labels.Labels) (*atomic.Int64, bool) { +func (s *activeSeriesStripe) findOrCreateEntryForSeries(fingerprint uint64, key uint32, series labels.Labels, nowNanos int64, nativeHistogram bool, labelsCopy func(labels.Labels) labels.Labels, tokens []uint32, ownedPositions []bool) (nanos *uatomic.Int64, entryTimeSet, created, owned bool) { s.mu.Lock() defer s.mu.Unlock() // Check if already exists within the entries. for ix, entry := range s.refs[fingerprint] { if labels.Equal(entry.lbs, series) { - return s.refs[fingerprint][ix].nanos, false + return s.refs[fingerprint][ix].nanos, false, false, entry.owned } } + // Ownership decides which counters this series contributes to, not whether it + // is tracked. A series this instance does not own is still an active series. + owned = isOwned(key, tokens, ownedPositions) + s.active++ + if owned { + s.owned++ + } if nativeHistogram { s.activeNativeHistogram++ } + e := activeSeriesEntry{ lbs: labelsCopy(series), - nanos: atomic.NewInt64(nowNanos), + key: key, + owned: owned, + nanos: uatomic.NewInt64(nowNanos), isNativeHistogram: nativeHistogram, } s.refs[fingerprint] = append(s.refs[fingerprint], e) - return e.nanos, true + return e.nanos, true, true, owned } -// nolint // Linter reports that this method is unused, but it is. -func (s *activeSeriesStripe) clear() { +// updateMetrics recounts this stripe's active and owned series, refreshing each +// entry's cached ownership if the ring changed. +// +// It removes nothing: releasing entries is Purge's job, and keeping idle entries +// is what allows owned to track what is in the head rather than what is in the +// idle window. +// +// There is no shortcut for "nothing expired" the way purge has, because entries +// leave the active window silently now that they are not deleted, so the counts +// have to be recomputed. The scan is the same order of work the purge it replaces +// performed on a tenant that was churning. +func (s *activeSeriesStripe) updateMetrics(keepUntil time.Time, tokensChanged bool, tokens []uint32, ownedPositions []bool) (active, owned, activeNativeHistogram int) { + keepUntilNanos := keepUntil.UnixNano() + s.mu.Lock() defer s.mu.Unlock() - s.oldestEntryTs.Store(0) - s.refs = map[uint64][]activeSeriesEntry{} - s.active = 0 + for _, entries := range s.refs { + for i := range entries { + if tokensChanged { + entries[i].owned = isOwned(entries[i].key, tokens, ownedPositions) + } + + if entries[i].owned { + owned++ + } + + if entries[i].nanos.Load() >= keepUntilNanos { + active++ + if entries[i].isNativeHistogram { + activeNativeHistogram++ + } + } + } + } + + s.active = active + s.owned = owned + s.activeNativeHistogram = activeNativeHistogram + + return active, owned, activeNativeHistogram } -func (s *activeSeriesStripe) purge(keepUntil time.Time) { +// purge removes entries last updated before keepUntil and returns the resulting +// counts for this stripe. +func (s *activeSeriesStripe) purge(keepUntil time.Time) (active, owned, activeNativeHistogram int) { keepUntilNanos := keepUntil.UnixNano() if oldest := s.oldestEntryTs.Load(); oldest > 0 && keepUntilNanos <= oldest { - // Nothing to do. - return + // Nothing to remove, so the counts cannot have changed. + s.mu.RLock() + defer s.mu.RUnlock() + + return s.active, s.owned, s.activeNativeHistogram } s.mu.Lock() defer s.mu.Unlock() - active := 0 - activeNativeHistogram := 0 - oldest := int64(math.MaxInt64) for fp, entries := range s.refs { - // Since we do expect very few fingerprint collisions, we - // have an optimized implementation for the common case. if len(entries) == 1 { ts := entries[0].nanos.Load() if ts < keepUntilNanos { @@ -188,6 +390,9 @@ func (s *activeSeriesStripe) purge(keepUntil time.Time) { } active++ + if entries[0].owned { + owned++ + } if entries[0].isNativeHistogram { activeNativeHistogram++ } @@ -197,8 +402,6 @@ func (s *activeSeriesStripe) purge(keepUntil time.Time) { continue } - // We have more entries, which means there's a collision, - // so we have to iterate over the entries. for i := 0; i < len(entries); { ts := entries[i].nanos.Load() if ts < keepUntilNanos { @@ -207,21 +410,20 @@ func (s *activeSeriesStripe) purge(keepUntil time.Time) { if ts < oldest { oldest = ts } - + active++ + if entries[i].owned { + owned++ + } + if entries[i].isNativeHistogram { + activeNativeHistogram++ + } i++ } } - // Either update or delete the entries in the map if cnt := len(entries); cnt == 0 { delete(s.refs, fp) } else { - active += cnt - for _, e := range entries { - if e.isNativeHistogram { - activeNativeHistogram++ - } - } s.refs[fp] = entries } } @@ -232,21 +434,42 @@ func (s *activeSeriesStripe) purge(keepUntil time.Time) { s.oldestEntryTs.Store(oldest) } s.active = active + s.owned = owned s.activeNativeHistogram = activeNativeHistogram + + return active, owned, activeNativeHistogram } -func (s *activeSeriesStripe) getActive() int { - s.mu.RLock() - defer s.mu.RUnlock() +// nolint // Linter reports that this method is unused, but it is. +func (s *activeSeriesStripe) clear() { + s.mu.Lock() + defer s.mu.Unlock() - return s.active + s.oldestEntryTs.Store(0) + s.refs = map[uint64][]activeSeriesEntry{} + s.active = 0 + s.owned = 0 + s.activeNativeHistogram = 0 } -func (s *activeSeriesStripe) getActiveNativeHistogram() int { - s.mu.RLock() - defer s.mu.RUnlock() +// isOwned reports whether the series with the given ring token is owned by this +// instance, meaning this instance is one of the replicas the ring selects for it. +// +// The answer is a binary search over the ring's token list followed by an array +// index into a bitmap precomputed once per ring change, so this is cheap enough +// to call on the push path. +// +// When ownership is unknown, because the ring has not been read yet or the bitmap +// does not match the token list, this returns true. Ownership is used to decide +// what counts against a limit, so the safe direction is to assume the series is +// ours: over-counting delays an unrelated scale-up, while under-counting would +// let a tenant exceed its limit. +func isOwned(key uint32, tokens []uint32, ownedPositions []bool) bool { + if len(tokens) == 0 || len(ownedPositions) != len(tokens) { + return true + } - return s.activeNativeHistogram + return ownedPositions[ring.SearchToken(tokens, key)] } // matchesAll returns true if the labels satisfy all given matchers. diff --git a/pkg/ingester/active_series_test.go b/pkg/ingester/active_series_test.go index 1f0d73bfa15..9d254ca2538 100644 --- a/pkg/ingester/active_series_test.go +++ b/pkg/ingester/active_series_test.go @@ -2,7 +2,9 @@ package ingester import ( "bytes" + "encoding/binary" "fmt" + "hash/fnv" "math" "strconv" "sync" @@ -12,6 +14,7 @@ import ( "github.com/prometheus/prometheus/model/labels" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) func copyFn(l labels.Labels) labels.Labels { return l } @@ -20,6 +23,15 @@ func fromLabelToLabels(ls []labels.Label) labels.Labels { return *(*labels.Labels)(unsafe.Pointer(&ls)) } +// noRingToken is the ring token passed by tests which are not exercising +// ownership. With no ring loaded, ownership is unknown and every series counts as +// owned, so these tests observe exactly the behaviour that predates owned-series +// tracking. +const noRingToken = uint32(0) + +// --- Active series behaviour. These tests predate owned-series tracking and are +// --- kept unchanged so that they continue to pin the active count's behaviour. + func TestActiveSeries_UpdateSeries(t *testing.T) { ls1 := []labels.Label{{Name: "a", Value: "1"}} ls2 := []labels.Label{{Name: "a", Value: "2"}} @@ -29,15 +41,15 @@ func TestActiveSeries_UpdateSeries(t *testing.T) { assert.Equal(t, 0, c.ActiveNativeHistogram()) labels1Hash := fromLabelToLabels(ls1).Hash() labels2Hash := fromLabelToLabels(ls2).Hash() - c.UpdateSeries(fromLabelToLabels(ls1), labels1Hash, time.Now(), true, copyFn) + c.UpdateSeries(fromLabelToLabels(ls1), labels1Hash, noRingToken, time.Now(), true, copyFn) assert.Equal(t, 1, c.Active()) assert.Equal(t, 1, c.ActiveNativeHistogram()) - c.UpdateSeries(fromLabelToLabels(ls1), labels1Hash, time.Now(), true, copyFn) + c.UpdateSeries(fromLabelToLabels(ls1), labels1Hash, noRingToken, time.Now(), true, copyFn) assert.Equal(t, 1, c.Active()) assert.Equal(t, 1, c.ActiveNativeHistogram()) - c.UpdateSeries(fromLabelToLabels(ls2), labels2Hash, time.Now(), true, copyFn) + c.UpdateSeries(fromLabelToLabels(ls2), labels2Hash, noRingToken, time.Now(), true, copyFn) assert.Equal(t, 2, c.Active()) assert.Equal(t, 2, c.ActiveNativeHistogram()) } @@ -56,7 +68,7 @@ func TestActiveSeries_Purge(t *testing.T) { c := NewActiveSeries() for i := range series { - c.UpdateSeries(fromLabelToLabels(series[i]), fromLabelToLabels(series[i]).Hash(), time.Unix(int64(i), 0), true, copyFn) + c.UpdateSeries(fromLabelToLabels(series[i]), fromLabelToLabels(series[i]).Hash(), noRingToken, time.Unix(int64(i), 0), true, copyFn) } c.Purge(time.Unix(int64(ttl+1), 0)) @@ -76,28 +88,418 @@ func TestActiveSeries_PurgeOpt(t *testing.T) { c := NewActiveSeries() now := time.Now() - c.UpdateSeries(ls1, ls1.Hash(), now.Add(-2*time.Minute), true, copyFn) - c.UpdateSeries(ls2, ls2.Hash(), now, true, copyFn) + c.UpdateSeries(ls1, ls1.Hash(), noRingToken, now.Add(-2*time.Minute), true, copyFn) + c.UpdateSeries(ls2, ls2.Hash(), noRingToken, now, true, copyFn) c.Purge(now) assert.Equal(t, 1, c.Active()) assert.Equal(t, 1, c.ActiveNativeHistogram()) - c.UpdateSeries(ls1, ls1.Hash(), now.Add(-1*time.Minute), true, copyFn) - c.UpdateSeries(ls2, ls2.Hash(), now, true, copyFn) + c.UpdateSeries(ls1, ls1.Hash(), noRingToken, now.Add(-1*time.Minute), true, copyFn) + c.UpdateSeries(ls2, ls2.Hash(), noRingToken, now, true, copyFn) c.Purge(now) assert.Equal(t, 1, c.Active()) assert.Equal(t, 1, c.ActiveNativeHistogram()) // This will *not* update the series, since there is already newer timestamp. - c.UpdateSeries(ls2, ls2.Hash(), now.Add(-1*time.Minute), true, copyFn) + c.UpdateSeries(ls2, ls2.Hash(), noRingToken, now.Add(-1*time.Minute), true, copyFn) c.Purge(now) assert.Equal(t, 1, c.Active()) assert.Equal(t, 1, c.ActiveNativeHistogram()) } +// --- Ownership helpers. + +// ownedPositionsFor builds an ownership bitmap over ringTokens from the subset of +// ring tokens this instance is a replica for. It lets these tests express +// ownership in terms of tokens, which reads more naturally, while the production +// code consumes the bitmap that ring.OwnedTokenPositions produces. +func ownedPositionsFor(ringTokens []uint32, ownedTokens ...uint32) []bool { + ownedSet := make(map[uint32]struct{}, len(ownedTokens)) + for _, token := range ownedTokens { + ownedSet[token] = struct{}{} + } + + positions := make([]bool, len(ringTokens)) + for position, token := range ringTokens { + _, positions[position] = ownedSet[token] + } + + return positions +} + +// testRingFingerprint derives a fingerprint from ring state the way the lifecycler +// does, so that tests passing identical state are seen as unchanged and tests +// passing different state are seen as a change. +func testRingFingerprint(ringTokens []uint32, ownedPositions []bool) uint64 { + h := fnv.New64a() + + var buf [4]byte + for _, token := range ringTokens { + binary.LittleEndian.PutUint32(buf[:], token) + _, _ = h.Write(buf[:]) + } + for _, owned := range ownedPositions { + if owned { + _, _ = h.Write([]byte{1}) + } else { + _, _ = h.Write([]byte{0}) + } + } + + return h.Sum64() +} + +// setRingState installs ring ownership state on c, reporting whether ownership +// changed. +func setRingState(c *ActiveSeries, ringTokens []uint32, ownedTokens ...uint32) bool { + positions := ownedPositionsFor(ringTokens, ownedTokens...) + return c.updateTokens(ringTokens, positions, testRingFingerprint(ringTokens, positions)) +} + +// updateMetricsWithRing calls UpdateMetrics with ownership expressed as the subset +// of ring tokens this instance is a replica for. +func updateMetricsWithRing(c *ActiveSeries, keepUntil time.Time, ringTokens []uint32, ownedTokens ...uint32) { + positions := ownedPositionsFor(ringTokens, ownedTokens...) + c.UpdateMetrics(keepUntil, ringTokens, positions, testRingFingerprint(ringTokens, positions)) +} + +// sumStripes totals the per-stripe counters, which are the source the cached +// totals are derived from. +func sumStripes(c *ActiveSeries) (active, owned, activeNativeHistogram int) { + for s := range numActiveSeriesStripes { + stripe := &c.stripes[s] + stripe.mu.RLock() + active += stripe.active + owned += stripe.owned + activeNativeHistogram += stripe.activeNativeHistogram + stripe.mu.RUnlock() + } + return active, owned, activeNativeHistogram +} + +func TestIsOwned(t *testing.T) { + // Ring with 4 tokens across 2 ingesters. With a replication factor of 1 each + // token range has a single owner, which is what makes the expectations below + // a simple alternation. + ringTokens := []uint32{100, 200, 300, 400} + ingester0 := ownedPositionsFor(ringTokens, 100, 300) + ingester1 := ownedPositionsFor(ringTokens, 200, 400) + + tests := []struct { + name string + key uint32 + ownedPositions []bool + expected bool + }{ + // Hash 50 → SearchToken finds 100 → ingester-0 owns it + {"hash 50 owned by ingester-0", 50, ingester0, true}, + {"hash 50 not owned by ingester-1", 50, ingester1, false}, + // Hash 150 → SearchToken finds 200 → ingester-1 owns it + {"hash 150 owned by ingester-1", 150, ingester1, true}, + {"hash 150 not owned by ingester-0", 150, ingester0, false}, + // Hash 250 → SearchToken finds 300 → ingester-0 owns it + {"hash 250 owned by ingester-0", 250, ingester0, true}, + {"hash 250 not owned by ingester-1", 250, ingester1, false}, + // Hash 350 → SearchToken finds 400 → ingester-1 owns it + {"hash 350 owned by ingester-1", 350, ingester1, true}, + {"hash 350 not owned by ingester-0", 350, ingester0, false}, + // Hash 450 → wraps around → SearchToken finds 100 → ingester-0 owns it + {"hash 450 wraps to ingester-0", 450, ingester0, true}, + {"hash 450 wraps, not ingester-1", 450, ingester1, false}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + assert.Equal(t, tc.expected, isOwned(tc.key, ringTokens, tc.ownedPositions)) + }) + } +} + +// TestIsOwned_UnknownOwnership covers the inputs which previously indexed a token +// slice without checking its length, which panicked on an empty ring. +func TestIsOwned_UnknownOwnership(t *testing.T) { + assert.True(t, isOwned(50, nil, nil), "empty ring must not panic and must not under-count") + assert.True(t, isOwned(50, []uint32{}, []bool{}), "empty ring must not panic and must not under-count") + assert.True(t, isOwned(50, []uint32{100, 200}, nil), "missing bitmap must not under-count") + assert.True(t, isOwned(50, []uint32{100, 200}, []bool{true}), "mismatched bitmap must not under-count") +} + +// TestActiveSeries_ActiveIsUnaffectedByOwnership is the regression test for the +// property that owned-series tracking must not disturb the pre-existing +// cortex_ingester_active_series gauge. The same series are pushed into two +// trackers, one with a ring loaded where half the tokens belong elsewhere and one +// with no ring at all, and the active counts must agree exactly. +func TestActiveSeries_ActiveIsUnaffectedByOwnership(t *testing.T) { + now := time.Now() + ringTokens := []uint32{100, 200} + + withRing := NewActiveSeries() + setRingState(withRing, ringTokens, 100) // owns token 100 only + + withoutRing := NewActiveSeries() + + // Keys alternate between the owned and the unowned token range. + for i := range 20 { + lbls := labels.FromStrings("__name__", "metric", "i", strconv.Itoa(i)) + key := uint32(50) + if i%2 == 1 { + key = 150 + } + + withRing.UpdateSeries(lbls, lbls.Hash(), key, now, i%3 == 0, copyFn) + withoutRing.UpdateSeries(lbls, lbls.Hash(), key, now, i%3 == 0, copyFn) + } + + assert.Equal(t, withoutRing.Active(), withRing.Active(), + "ownership tracking must not change the active series count") + assert.Equal(t, withoutRing.ActiveNativeHistogram(), withRing.ActiveNativeHistogram(), + "ownership tracking must not change the active native histogram count") + assert.Equal(t, 20, withRing.Active()) + + // Ownership is the only thing that differs between the two. + assert.Equal(t, 10, withRing.Owned(), "half the keys fall in a token range owned elsewhere") + assert.Equal(t, 20, withoutRing.Owned(), "with no ring, ownership is unknown and everything counts") +} + +func TestActiveSeries_OwnedCount_NoRingLoaded(t *testing.T) { + // Before the ring is read, ownership is unknown, so Owned() equals Active(). + c := NewActiveSeries() + now := time.Now() + + lbls1 := labels.FromStrings("__name__", "metric_1", "job", "test") + lbls2 := labels.FromStrings("__name__", "metric_2", "job", "test") + + c.UpdateSeries(lbls1, lbls1.Hash(), noRingToken, now, false, copyFn) + c.UpdateSeries(lbls2, lbls2.Hash(), noRingToken, now, false, copyFn) + + assert.Equal(t, 2, c.Active()) + assert.Equal(t, 2, c.Owned()) +} + +func TestActiveSeries_OwnedCount_WithRingLoaded(t *testing.T) { + // With the ring loaded, an unowned series is still tracked as active but does + // not count towards owned. + c := NewActiveSeries() + now := time.Now() + + ringTokens := []uint32{100, 200} + setRingState(c, ringTokens, 100) // owns token 100 + + // key=50 → SearchToken finds 100 → owned + lblsOwned := labels.FromStrings("__name__", "metric_owned", "job", "test") + c.UpdateSeries(lblsOwned, lblsOwned.Hash(), 50, now, false, copyFn) + + // key=150 → SearchToken finds 200 → not owned + lblsUnowned := labels.FromStrings("__name__", "metric_not_owned", "job", "test") + c.UpdateSeries(lblsUnowned, lblsUnowned.Hash(), 150, now, false, copyFn) + + assert.Equal(t, 2, c.Active(), "both series are active regardless of ownership") + assert.Equal(t, 1, c.Owned(), "only the owned series counts towards owned") +} + +// TestActiveSeries_OwnedExceedsActiveForIdleSeries covers the behaviour the whole +// feature rests on. Entries are retained until Purge is called with the head's +// minimum time, so a series which has gone idle still counts towards owned while +// it is still held in the head. owned therefore tracks what is in memory, which is +// what the series limit is protecting, rather than what is in the idle window. +func TestActiveSeries_OwnedExceedsActiveForIdleSeries(t *testing.T) { + c := NewActiveSeries() + now := time.Now() + idleCutoff := now.Add(-10 * time.Minute) + + ringTokens := []uint32{100} + setRingState(c, ringTokens, 100) // owns everything + + // One recent series and three which have gone idle. + recent := labels.FromStrings("__name__", "recent") + c.UpdateSeries(recent, recent.Hash(), 50, now, false, copyFn) + for i := range 3 { + lbls := labels.FromStrings("__name__", "idle", "i", strconv.Itoa(i)) + c.UpdateSeries(lbls, lbls.Hash(), 50, now.Add(-time.Hour), false, copyFn) + } + + updateMetricsWithRing(c, idleCutoff, ringTokens, 100) + + assert.Equal(t, 1, c.Active(), "only the recent series is inside the idle window") + assert.Equal(t, 4, c.Owned(), "all four are still held, so all four count towards owned") + assert.Greater(t, c.Owned(), c.Active(), "owned exceeds active for an idle-heavy tenant") +} + +// TestActiveSeries_UpdateMetricsRetainsExpiredEntries pins down that the periodic +// cycle removes nothing, which is what allows owned to outlive the idle window. +func TestActiveSeries_UpdateMetricsRetainsExpiredEntries(t *testing.T) { + c := NewActiveSeries() + now := time.Now() + idleCutoff := now.Add(-30 * time.Minute) + + ringTokens := []uint32{100} + setRingState(c, ringTokens, 100) + + old := labels.FromStrings("__name__", "old_metric") + c.UpdateSeries(old, old.Hash(), 50, now.Add(-time.Hour), false, copyFn) + recent := labels.FromStrings("__name__", "recent_metric") + c.UpdateSeries(recent, recent.Hash(), 50, now, false, copyFn) + + updateMetricsWithRing(c, idleCutoff, ringTokens, 100) + + assert.Equal(t, 1, c.Active()) + assert.Equal(t, 2, c.Owned(), "the expired entry is retained and still owned") + + // Repeated cycles must be stable rather than progressively dropping entries. + updateMetricsWithRing(c, idleCutoff, ringTokens, 100) + assert.Equal(t, 1, c.Active()) + assert.Equal(t, 2, c.Owned()) + + // Only a purge releases it. + c.Purge(idleCutoff) + assert.Equal(t, 1, c.Active()) + assert.Equal(t, 1, c.Owned()) +} + +// TestActiveSeries_RingChangeMovesOwnedNotActive checks that losing ownership of a +// series changes only the owned count. The series is still held in this ingester's +// head, so it must remain active and must not be deleted. +func TestActiveSeries_RingChangeMovesOwnedNotActive(t *testing.T) { + c := NewActiveSeries() + now := time.Now() + idleCutoff := now.Add(-10 * time.Minute) + + setRingState(c, []uint32{100, 200}, 100) + + lbls := labels.FromStrings("__name__", "metric_1", "job", "test") + c.UpdateSeries(lbls, lbls.Hash(), 50, now, false, copyFn) + + assert.Equal(t, 1, c.Active()) + assert.Equal(t, 1, c.Owned()) + + // The ring changes and token 100's range now belongs elsewhere. + updateMetricsWithRing(c, idleCutoff, []uint32{100, 200, 300}, 200) + + assert.Equal(t, 1, c.Active(), "the series is still held here, so it is still active") + assert.Equal(t, 0, c.Owned(), "but it is no longer owned") + + // And ownership can come back without the series having to be re-pushed. + updateMetricsWithRing(c, idleCutoff, []uint32{100, 200, 300}, 100) + assert.Equal(t, 1, c.Active()) + assert.Equal(t, 1, c.Owned()) +} + +// TestActiveSeries_CachedTotalsMatchStripes guards the cached totals that +// PreCreation reads. They are incremented as series are created and recomputed by +// the periodic cycle, so a mismatch would mean limits are enforced against a +// number that has drifted from reality. +func TestActiveSeries_CachedTotalsMatchStripes(t *testing.T) { + c := NewActiveSeries() + now := time.Now() + ringTokens := []uint32{100, 200} + setRingState(c, ringTokens, 100) + + for i := range 50 { + lbls := labels.FromStrings("__name__", "metric", "i", strconv.Itoa(i)) + key := uint32(50) + if i%2 == 1 { + key = 150 + } + c.UpdateSeries(lbls, lbls.Hash(), key, now, i%5 == 0, copyFn) + } + + assertTotalsMatchStripes := func(stage string) { + t.Helper() + active, owned, activeNativeHistogram := sumStripes(c) + assert.Equal(t, active, c.Active(), "active total drifted from stripes after %s", stage) + assert.Equal(t, owned, c.Owned(), "owned total drifted from stripes after %s", stage) + assert.Equal(t, activeNativeHistogram, c.ActiveNativeHistogram(), "native histogram total drifted from stripes after %s", stage) + } + + assertTotalsMatchStripes("creation") + + updateMetricsWithRing(c, now.Add(-time.Minute), ringTokens, 100) + assertTotalsMatchStripes("periodic cycle") + + c.Purge(now.Add(-time.Minute)) + assertTotalsMatchStripes("purge") + + c.clear() + assertTotalsMatchStripes("clear") + assert.Equal(t, 0, c.Active()) + assert.Equal(t, 0, c.Owned()) +} + +func TestActiveSeries_NativeHistogram_Owned(t *testing.T) { + c := NewActiveSeries() + now := time.Now() + + ringTokens := []uint32{100} + setRingState(c, ringTokens, 100) + + lbls := labels.FromStrings("__name__", "histogram_metric", "job", "test") + c.UpdateSeries(lbls, lbls.Hash(), 50, now, true, copyFn) + + assert.Equal(t, 1, c.Active()) + assert.Equal(t, 1, c.Owned()) + assert.Equal(t, 1, c.ActiveNativeHistogram()) +} + +func TestActiveSeries_ExistingSeriesKeepsOwnership(t *testing.T) { + // A second sample for a series already tracked must not duplicate it, and must + // not disturb its ownership. + c := NewActiveSeries() + now := time.Now() + later := now.Add(1 * time.Minute) + + setRingState(c, []uint32{100, 200}, 100) + + lbls := labels.FromStrings("__name__", "metric_1", "job", "test") + c.UpdateSeries(lbls, lbls.Hash(), 50, now, false, copyFn) + require.Equal(t, 1, c.Active()) + require.Equal(t, 1, c.Owned()) + + c.UpdateSeries(lbls, lbls.Hash(), 50, later, false, copyFn) + assert.Equal(t, 1, c.Active()) + assert.Equal(t, 1, c.Owned()) +} + +func TestUpdateTokens_DetectsChange(t *testing.T) { + c := NewActiveSeries() + + // First call should detect change (from empty to something) + assert.True(t, setRingState(c, []uint32{100, 200}, 100)) + + // Same state again — no change + assert.False(t, setRingState(c, []uint32{100, 200}, 100)) + + // Different ring tokens — change detected + assert.True(t, setRingState(c, []uint32{100, 200, 300}, 100)) + + // Same tokens but different ownership — must still be detected, because a peer + // changing state alters the replica set without moving any token. + assert.True(t, setRingState(c, []uint32{100, 200, 300}, 100, 200)) +} + +func TestActiveSeries_UpdateTokens_ImmutableSnapshots(t *testing.T) { + // Verify that updateTokens publishes a new ringState each time rather than + // mutating the previous one, which is what makes the atomic.Pointer safe. + c := NewActiveSeries() + + setRingState(c, []uint32{100, 200}, 100) + state1 := c.ring.Load() + + setRingState(c, []uint32{100, 200, 300}, 100, 300) + state2 := c.ring.Load() + + assert.NotEqual(t, state1, state2) + assert.Len(t, state1.tokens, 2) + assert.Len(t, state2.tokens, 3) + assert.Len(t, state1.ownedPositions, 2) + assert.Len(t, state2.ownedPositions, 3) +} + +// --- Benchmarks. These predate owned-series tracking and are kept so that the +// --- push and purge paths stay comparable against earlier numbers. + var activeSeriesTestGoroutines = []int{50, 100, 500} func BenchmarkActiveSeriesTest_single_series(b *testing.B) { @@ -125,7 +527,7 @@ func benchmarkActiveSeriesConcurrencySingleSeries(b *testing.B, goroutines int) for ix := range max { now = now.Add(time.Duration(ix) * time.Millisecond) - c.UpdateSeries(series, labelhash, now, false, copyFn) + c.UpdateSeries(series, labelhash, noRingToken, now, false, copyFn) } }) } @@ -145,17 +547,21 @@ func BenchmarkActiveSeries_UpdateSeries(b *testing.B) { } name := nameBuf.String() + // NOTE: this deliberately does not use b.Loop(). A benchmark may only run one + // b.Loop() loop, and sizing the series slice requires knowing the iteration + // count before the loop starts, so the b.N form is the correct one here. series := make([]labels.Labels, b.N) labelhash := make([]uint64, b.N) - for s := 0; b.Loop(); s++ { + for s := 0; s < b.N; s++ { series[s] = labels.FromStrings(name, name+strconv.Itoa(s)) labelhash[s] = series[s].Hash() } now := time.Now().UnixNano() + b.ResetTimer() - for ix := 0; b.Loop(); ix++ { - c.UpdateSeries(series[ix], labelhash[ix], time.Unix(0, now+int64(ix)), false, copyFn) + for ix := 0; ix < b.N; ix++ { + c.UpdateSeries(series[ix], labelhash[ix], noRingToken, time.Unix(0, now+int64(ix)), false, copyFn) } } @@ -187,9 +593,9 @@ func benchmarkPurge(b *testing.B, twice bool) { // Prepare series for ix, s := range series { if ix < numExpiresSeries { - c.UpdateSeries(s, labelhash[ix], now.Add(-time.Minute), false, copyFn) + c.UpdateSeries(s, labelhash[ix], noRingToken, now.Add(-time.Minute), false, copyFn) } else { - c.UpdateSeries(s, labelhash[ix], now, false, copyFn) + c.UpdateSeries(s, labelhash[ix], noRingToken, now, false, copyFn) } } @@ -206,3 +612,93 @@ func benchmarkPurge(b *testing.B, twice bool) { } } } + +// --- Benchmarks for the owned-series paths. + +// BenchmarkActiveSeries_UpdateSeries_Owned measures the push path with a ring +// loaded, which is the cost ownership tracking adds per new series: one binary +// search over the ring's tokens plus one array index. +func BenchmarkActiveSeries_UpdateSeries_Owned(b *testing.B) { + const numRingTokens = 100 * 512 + + ringTokens := make([]uint32, numRingTokens) + for i := range ringTokens { + ringTokens[i] = uint32(i) * 128 + } + ownedPositions := make([]bool, numRingTokens) + for i := range ownedPositions { + ownedPositions[i] = i%3 == 0 + } + + for _, withRing := range []bool{false, true} { + name := "ring_not_loaded" + if withRing { + name = "ring_loaded" + } + + b.Run(name, func(b *testing.B) { + c := NewActiveSeries() + if withRing { + c.updateTokens(ringTokens, ownedPositions, 1) + } + + series := make([]labels.Labels, b.N) + labelhash := make([]uint64, b.N) + for s := 0; s < b.N; s++ { + series[s] = labels.FromStrings("__name__", "metric", "i", strconv.Itoa(s)) + labelhash[s] = series[s].Hash() + } + + now := time.Now() + b.ReportAllocs() + b.ResetTimer() + + for ix := 0; ix < b.N; ix++ { + c.UpdateSeries(series[ix], labelhash[ix], uint32(ix)*7919, now, false, copyFn) + } + }) + } +} + +// BenchmarkActiveSeries_UpdateMetrics measures the periodic recount, separating the +// common case where the ring has not changed from the case where every entry's +// ownership has to be re-evaluated. +func BenchmarkActiveSeries_UpdateMetrics(b *testing.B) { + const ( + numSeries = 100000 + numRingTokens = 1000 + ) + + ringTokens := make([]uint32, numRingTokens) + for i := range ringTokens { + ringTokens[i] = uint32(i) * 4096 + } + ownedPositions := make([]bool, numRingTokens) + for i := range ownedPositions { + ownedPositions[i] = i%3 == 0 + } + + now := time.Now() + c := NewActiveSeries() + c.updateTokens(ringTokens, ownedPositions, 1) + for i := range numSeries { + lbls := labels.FromStrings("__name__", "metric", "i", strconv.Itoa(i)) + c.UpdateSeries(lbls, lbls.Hash(), uint32(i)*7919, now, false, copyFn) + } + + b.Run("ring_unchanged", func(b *testing.B) { + b.ReportAllocs() + for b.Loop() { + c.UpdateMetrics(now.Add(-time.Minute), ringTokens, ownedPositions, 1) + } + }) + + b.Run("ring_changed", func(b *testing.B) { + b.ReportAllocs() + fingerprint := uint64(1) + for b.Loop() { + fingerprint++ + c.UpdateMetrics(now.Add(-time.Minute), ringTokens, ownedPositions, fingerprint) + } + }) +} diff --git a/pkg/ingester/active_series_tracker_test.go b/pkg/ingester/active_series_tracker_test.go index adcd2a03309..69aa18e01e0 100644 --- a/pkg/ingester/active_series_tracker_test.go +++ b/pkg/ingester/active_series_tracker_test.go @@ -200,7 +200,7 @@ func TestActiveSeries_Basic(t *testing.T) { now := time.Now() s := labels.FromStrings("__name__", "test", "job", "app") - as.UpdateSeries(s, s.Hash(), now, false, func(l labels.Labels) labels.Labels { return l.Copy() }) + as.UpdateSeries(s, s.Hash(), 0, now, false, func(l labels.Labels) labels.Labels { return l.Copy() }) assert.Equal(t, 1, as.Active()) } diff --git a/pkg/ingester/client/compat.go b/pkg/ingester/client/compat.go index 20f745c52eb..377dbbb8d62 100644 --- a/pkg/ingester/client/compat.go +++ b/pkg/ingester/client/compat.go @@ -10,6 +10,7 @@ import ( storecache "github.com/thanos-io/thanos/pkg/store/cache" "github.com/cortexproject/cortex/pkg/cortexpb" + "github.com/cortexproject/cortex/pkg/util" ) // ToQueryRequest builds a QueryRequest proto. @@ -294,10 +295,10 @@ func FastFingerprint(ls []cortexpb.LabelAdapter) model.Fingerprint { var result uint64 for _, l := range ls { - sum := hashNew() - sum = hashAdd(sum, l.Name) - sum = hashAddByte(sum, model.SeparatorByte) - sum = hashAdd(sum, l.Value) + sum := util.HashNew() + sum = util.HashAdd(sum, l.Name) + sum = util.HashAddByte(sum, model.SeparatorByte) + sum = util.HashAdd(sum, l.Value) result ^= sum } return model.Fingerprint(result) diff --git a/pkg/ingester/ingester.go b/pkg/ingester/ingester.go index 3a6c03449c6..fdc04465a93 100644 --- a/pkg/ingester/ingester.go +++ b/pkg/ingester/ingester.go @@ -142,6 +142,19 @@ type Config struct { HeadQueriedSeriesMetricsSampleRate float64 `yaml:"head_queried_series_metrics_sample_rate"` HeadQueriedSeriesMetricsWindows cortex_tsdb.DurationList `yaml:"head_queried_series_metrics_windows"` + // OwnedSeriesMetricsEnabled enables tracking of owned series per user. + // When enabled, the ingester computes series ownership based on the ring and + // emits the cortex_ingester_owned_series metric. Does NOT change limit enforcement + // behavior — use OwnedSeriesLimitEnforcementEnabled for that. + OwnedSeriesMetricsEnabled bool `yaml:"owned_series_metrics_enabled"` + + // OwnedSeriesLimitEnforcementEnabled enables using owned series count for limit + // enforcement in PreCreation(). When enabled (requires OwnedSeriesMetricsEnabled=true), + // both the per-user series limit and instance-level max_series limit use owned count + // instead of Head().NumSeries(), preventing false throttling during resharding. + // If OwnedSeriesMetricsEnabled is false, this flag has no effect (falls back to old behavior). + OwnedSeriesLimitEnforcementEnabled bool `yaml:"owned_series_limit_enforcement_enabled"` + // Use blocks storage. BlocksStorageConfig cortex_tsdb.BlocksStorageConfig `yaml:"-"` @@ -214,6 +227,9 @@ func (cfg *Config) RegisterFlags(f *flag.FlagSet) { cfg.HeadQueriedSeriesMetricsWindows = cortex_tsdb.DurationList{2 * time.Hour} f.Var(&cfg.HeadQueriedSeriesMetricsWindows, "ingester.head-queried-series-metrics-windows", "Time windows to expose head queried series metrics. Also controls how long per-metric-name cardinality is reported after last query.") + f.BoolVar(&cfg.OwnedSeriesMetricsEnabled, "ingester.owned-series-metrics-enabled", false, "Enable tracking of owned series per user. When enabled, the ingester computes series ownership based on the ring and emits cortex_ingester_owned_series metric.") + f.BoolVar(&cfg.OwnedSeriesLimitEnforcementEnabled, "ingester.owned-series-limit-enforcement-enabled", false, "Use owned series count for limit enforcement. Requires owned-series-metrics-enabled. When enabled, PreCreation uses owned count instead of Head().NumSeries() for both per-user and instance-level limits.") + f.BoolVar(&cfg.UploadCompactedBlocksEnabled, "ingester.upload-compacted-blocks-enabled", true, "Enable uploading compacted blocks.") f.StringVar(&cfg.IgnoreSeriesLimitForMetricNames, "ingester.ignore-series-limit-for-metric-names", "", "Comma-separated list of metric names, for which -ingester.max-series-per-metric and -ingester.max-global-series-per-metric limits will be ignored. Does not affect max-series-per-user or max-global-series-per-metric limits.") f.StringVar(&cfg.AdminLimitMessage, "ingester.admin-limit-message", "please contact administrator to raise it", "Customize the message contained in limit errors") @@ -389,6 +405,19 @@ type userTSDB struct { instanceSeriesCount *atomic.Int64 // Shared across all userTSDB instances created by ingester. instanceLimitsFn func() *InstanceLimits + // ownedSeriesLimitEnabled controls whether PreCreation uses activeSeries.Owned() + // for limit checks instead of Head().NumSeries(). Only true when BOTH + // OwnedSeriesMetricsEnabled AND OwnedSeriesLimitEnforcementEnabled are set. + ownedSeriesLimitEnabled bool + + // instanceOwnedCount tracks total owned series across all tenants on this ingester. + // Recalculated every updateActiveSeries cycle (1 min). Used for instance-level + // max_series limit when ownedSeriesLimitEnabled is true. + // NOTE: Up to 1 minute stale after ring changes. This is acceptable because + // staleness is conservative (overcounts) and this limit protects against OOM, + // not customer-facing throttle errors. + instanceOwnedCount *atomic.Int64 + stateMtx sync.RWMutex state tsdbState pushesInFlight sync.WaitGroup // Increased with stateMtx read lock held, only if state == active or activeShipping. @@ -517,22 +546,79 @@ func (u *userTSDB) compactHead(ctx context.Context, blockDuration int64) error { return u.db.CompactOOOHead(ctx) } +// purgeActiveSeriesToHead releases active-series entries whose series are no longer +// in the TSDB head, by purging everything last seen before the head's new minimum +// time. Called after a successful head compaction. +// +// This is the second tier of a two-tier retention scheme. The periodic cycle +// recounts without removing anything, so that an idle series still held in the head +// keeps counting towards owned; this tier removes an entry once the series it +// describes has actually left memory. +// +// A series appended immediately after the head's minimum time is read could have +// its entry dropped here while still being in the head, which would undercount +// owned. That self-corrects on the series' next sample, which recreates the entry. +func (u *userTSDB) purgeActiveSeriesToHead() { + h := u.Head() + + // An empty head means nothing is in memory, so every entry can go. MinTime is + // not usable in this case: Prometheus reports math.MaxInt64 for an empty head. + if h.NumSeries() == 0 { + u.activeSeries.clear() + return + } + + minTime := h.MinTime() + if minTime <= 0 || minTime == math.MaxInt64 { + return + } + + u.activeSeries.Purge(time.UnixMilli(minTime)) +} + // PreCreation implements SeriesLifecycleCallback interface. func (u *userTSDB) PreCreation(metric labels.Labels) error { if u.limiter == nil { return nil } - // Verify ingester's global limit + // Verify ingester's global limit (instance-level max_series). + // When limit enforcement is enabled, use instanceOwnedCount which reflects + // only series this ingester currently owns according to the ring. + // NOTE: instanceOwnedCount is recalculated every ~1 min in updateActiveSeries. + // Up to 1 min stale after ring changes, but conservative (overcounts). gl := u.instanceLimitsFn() if gl != nil && gl.MaxInMemorySeries > 0 { - if series := u.instanceSeriesCount.Load(); series >= gl.MaxInMemorySeries { + var instanceCount int64 + // A negative owned count means no cycle has computed one yet. Fall back to + // the total in-memory series count so that OOM protection still applies + // during the first cycle after startup. A count of zero is a real count and + // is used as-is. + if owned := u.instanceOwnedCount.Load(); u.ownedSeriesLimitEnabled && owned >= 0 { + instanceCount = owned + } else { + instanceCount = u.instanceSeriesCount.Load() + } + if instanceCount >= gl.MaxInMemorySeries { return errMaxSeriesLimitReached } } - // Total series limit. - if err := u.limiter.AssertMaxSeriesPerUser(u.userID, int(u.Head().NumSeries())); err != nil { + // Per-user series limit. + // When limit enforcement is enabled (flag 2 + flag 1), use activeSeries.Owned() + // which counts the series this ingester holds and owns, rather than everything + // left in the head. This prevents false throttling during scale-up, where + // Head().NumSeries() stays high for up to a head-compaction cycle while the + // local limit has already dropped because more ingesters joined the ring. + // + // Owned() needs no startup fallback: it is maintained as series are created, so + // it is accurate from the very first sample rather than only after the first + // periodic cycle. + seriesCount := int(u.Head().NumSeries()) + if u.ownedSeriesLimitEnabled { + seriesCount = u.activeSeries.Owned() + } + if err := u.limiter.AssertMaxSeriesPerUser(u.userID, seriesCount); err != nil { return err } @@ -733,6 +819,14 @@ type TSDBState struct { // Number of series in memory, across all tenants. seriesCount atomic.Int64 + // Number of owned series across all tenants. Recalculated every updateActiveSeries + // cycle (~1 min). Used for instance-level max_series when limit enforcement is enabled. + // + // Negative means "not computed yet", which is distinct from a genuine zero. The + // distinction matters: zero is a legitimate count that must be trusted, whereas + // before the first cycle there is no owned count to compare a limit against. + ownedSeriesCount atomic.Int64 + // Head compactions metrics. compactionsTriggered prometheus.Counter compactionsFailed prometheus.Counter @@ -764,7 +858,7 @@ func newTSDBState(bucketClient objstore.Bucket, registerer prometheus.Registerer idleTsdbChecks.WithLabelValues(string(tsdbTenantMarkedForDeletion)) idleTsdbChecks.WithLabelValues(string(tsdbIdleClosed)) - return TSDBState{ + state := TSDBState{ dbs: make(map[string]*userTSDB), bucket: bucketClient, tsdbMetrics: newTSDBMetrics(registerer), @@ -807,6 +901,13 @@ func newTSDBState(bucketClient objstore.Bucket, registerer prometheus.Registerer idleTsdbChecks: idleTsdbChecks, } + + // No owned count has been computed yet. Until the first updateActiveSeries + // cycle runs, the instance-level limit falls back to the total in-memory + // series count so that OOM protection is never switched off. + state.ownedSeriesCount.Store(-1) + + return state } // New returns a new Ingester that uses Cortex block storage instead of chunks storage. @@ -1183,13 +1284,29 @@ func (i *Ingester) getMaxExemplars(userID string) int64 { func (i *Ingester) updateActiveSeries(ctx context.Context) { purgeTime := time.Now().Add(-i.cfg.ActiveSeriesMetricsIdleTimeout) + // When owned metrics are enabled, recalculate the instance-level owned count + // from scratch each cycle. This avoids drift from edge cases (missed decrements + // during ring changes or tenant deletions). The loop already iterates all userTSDBs + // and calls Owned(), so this is essentially free (one int64 addition per tenant). + var totalOwnedCount int64 + for _, userID := range i.getTSDBUsers() { userDB, err := i.getTSDB(userID) if err != nil || userDB == nil { continue } - userDB.activeSeries.Purge(purgeTime) + if i.cfg.OwnedSeriesMetricsEnabled { + // Use UpdateMetrics which handles both purge AND ownership re-evaluation. + ringTokens, ownedPositions, ringFingerprint := i.lifecycler.GetOwnedTokenPositions() + userDB.activeSeries.UpdateMetrics(purgeTime, ringTokens, ownedPositions, ringFingerprint) + owned := userDB.activeSeries.Owned() + i.metrics.ownedSeriesPerUser.WithLabelValues(userID).Set(float64(owned)) + totalOwnedCount += int64(owned) + } else { + userDB.activeSeries.Purge(purgeTime) + } + i.metrics.activeSeriesPerUser.WithLabelValues(userID).Set(float64(userDB.activeSeries.Active())) i.metrics.activeNHSeriesPerUser.WithLabelValues(userID).Set(float64(userDB.activeSeries.ActiveNativeHistogram())) i.metrics.headMetricNamesPerUser.WithLabelValues(userID).Set(float64(userDB.seriesInMetric.ActiveMetricNames())) @@ -1202,6 +1319,11 @@ func (i *Ingester) updateActiveSeries(ctx context.Context) { userDB.trackerCounter.updateConfig(ctx, userDB.db, trackers) userDB.trackerCounter.updateMetrics(i.metrics.activeSeriesPerTracker, userID, trackers) } + + // Store the instance-level owned count for use in PreCreation's max_series check. + if i.cfg.OwnedSeriesMetricsEnabled { + i.TSDBState.ownedSeriesCount.Store(totalOwnedCount) + } } func (i *Ingester) updateActiveQueriedSeries(ctx context.Context) { @@ -1557,6 +1679,18 @@ func (i *Ingester) Push(ctx context.Context, req *cortexpb.WriteRequest) (*corte return nil, wrapWithUser(errors.Errorf("out-of-order label set found when push: %s", tsLabels), userID) } tsLabelsHash := tsLabels.Hash() + + // Compute ring token for this series (same hash the distributor uses for routing). + // Used by ActiveSeries to track ownership. When flag is off, tsToken=0 and + // ActiveSeries skips ownership checks. + var tsToken uint32 + if i.cfg.OwnedSeriesMetricsEnabled { + tsToken, err = ring.TokenForLabels(userID, ts.Labels, i.cfg.DistributorShardByAllLabels) + if err != nil { + return nil, wrapWithUser(err, userID) + } + } + ref, copiedLabels := app.GetRef(tsLabels, tsLabelsHash) // To find out if any sample was added to this series, we keep old value. @@ -1681,7 +1815,7 @@ func (i *Ingester) Push(ctx context.Context, req *cortexpb.WriteRequest) (*corte isNHAppended := succeededHistogramsCount > oldSucceededHistogramsCount shouldUpdateSeries := (succeededSamplesCount > oldSucceededSamplesCount) || isNHAppended if i.cfg.ActiveSeriesMetricsEnabled && shouldUpdateSeries { - db.activeSeries.UpdateSeries(tsLabels, tsLabelsHash, startAppend, isNHAppended, func(l labels.Labels) labels.Labels { + db.activeSeries.UpdateSeries(tsLabels, tsLabelsHash, tsToken, startAppend, isNHAppended, func(l labels.Labels) labels.Labels { // we must already have copied the labels if succeededSamplesCount or succeededHistogramsCount has been incremented. return copiedLabels }) @@ -3056,6 +3190,8 @@ func (i *Ingester) createTSDB(userID string) (*userTSDB, error) { instanceSeriesCount: &i.TSDBState.seriesCount, interner: util.NewLruInterner(i.cfg.LabelsStringInterningEnabled), labelsStringInterningEnabled: i.cfg.LabelsStringInterningEnabled, + ownedSeriesLimitEnabled: i.cfg.OwnedSeriesMetricsEnabled && i.cfg.OwnedSeriesLimitEnforcementEnabled, + instanceOwnedCount: &i.TSDBState.ownedSeriesCount, blockRetentionPeriod: i.cfg.BlocksStorageConfig.TSDB.Retention.Milliseconds(), postingCache: postingCache, @@ -3450,7 +3586,7 @@ func (i *Ingester) compactionLoop(ctx context.Context) error { } // Lets create the slot based on the hash id - i := int(client.HashAdd32(client.HashNew32(), i.lifecycler.ID) % 10) + i := int(util.HashAdd32(util.HashNew32(), i.lifecycler.ID) % 10) return i, 10 } ticker := util.NewSlottedTicker(infoFunc, i.cfg.BlocksStorageConfig.TSDB.HeadCompactionInterval, 1) @@ -3522,6 +3658,14 @@ func (i *Ingester) compactBlocks(ctx context.Context, force bool, allowed *users level.Warn(logutil.WithContext(ctx, i.logger)).Log("msg", "TSDB blocks compaction for user has failed", "user", userID, "err", err, "compactReason", reason) } else { level.Debug(logutil.WithContext(ctx, i.logger)).Log("msg", "TSDB blocks compaction completed successfully", "user", userID, "compactReason", reason) + + // Compaction has moved series out of the head, so the entries tracking + // them can be released. This is the only place entries are dropped when + // owned-series tracking is on, which is what keeps the owned count + // aligned with what is in memory rather than with the idle window. + if i.cfg.OwnedSeriesMetricsEnabled { + userDB.purgeActiveSeriesToHead() + } } return nil diff --git a/pkg/ingester/metrics.go b/pkg/ingester/metrics.go index 3ad21faad6d..54492ea9450 100644 --- a/pkg/ingester/metrics.go +++ b/pkg/ingester/metrics.go @@ -61,6 +61,7 @@ type ingesterMetrics struct { activeSeriesPerUser *prometheus.GaugeVec activeNHSeriesPerUser *prometheus.GaugeVec headMetricNamesPerUser *prometheus.GaugeVec + ownedSeriesPerUser *prometheus.GaugeVec activeQueriedSeriesPerUser *prometheus.GaugeVec headQueriedSeriesPerUser *prometheus.GaugeVec limitsPerLabelSet *prometheus.GaugeVec @@ -330,6 +331,11 @@ func newIngesterMetrics(r prometheus.Registerer, Help: "Number of unique metric names in the TSDB head per user.", }, []string{"user"}), + ownedSeriesPerUser: prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Name: "cortex_ingester_owned_series", + Help: "Number of series this ingester currently owns per user according to the ring.", + }, []string{"user"}), + // Not registered automatically, but only if activeSeriesEnabled is true. activeSeriesPerTracker: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "cortex_ingester_active_series_per_tracker", @@ -388,6 +394,7 @@ func newIngesterMetrics(r prometheus.Registerer, r.MustRegister(m.activeSeriesPerUser) r.MustRegister(m.activeNHSeriesPerUser) r.MustRegister(m.headMetricNamesPerUser) + r.MustRegister(m.ownedSeriesPerUser) r.MustRegister(m.activeSeriesPerTracker) } diff --git a/pkg/ring/lifecycler.go b/pkg/ring/lifecycler.go index 1db33a929f0..a1fa7b894f2 100644 --- a/pkg/ring/lifecycler.go +++ b/pkg/ring/lifecycler.go @@ -157,6 +157,17 @@ type Lifecycler struct { zonesCount int zones []string + // Ownership of ring token positions by this instance. ownedTokens is the + // ring's full sorted token list and ownedPositions is parallel to it. + // + // These are recomputed only when the ring topology actually changes, tracked + // by ownedFingerprint. Once published they are never mutated in place, so + // readers may keep using a snapshot after releasing countersLock. + ownedTokens []uint32 + ownedPositions []bool + ownedFingerprint uint64 + ownedFingerprintOK bool + lifecyclerMetrics *LifecyclerMetrics logger log.Logger @@ -283,7 +294,7 @@ func (i *Lifecycler) CheckReady(ctx context.Context) error { func (i *Lifecycler) checkRingHealthForReadiness(ctx context.Context) error { // Ensure the instance holds some tokens. - if len(i.getTokens()) == 0 { + if len(i.GetTokens()) == 0 { return fmt.Errorf("this instance owns no tokens") } @@ -363,12 +374,35 @@ func (i *Lifecycler) ChangeState(ctx context.Context, state InstanceState) error return <-errCh } -func (i *Lifecycler) getTokens() Tokens { +func (i *Lifecycler) GetTokens() Tokens { i.stateMtx.RLock() defer i.stateMtx.RUnlock() return i.tokenFile.Tokens } +// GetOwnedTokenPositions returns the ring's full sorted token list together with a +// parallel bitmap of the token positions owned by this instance, meaning the +// positions whose replica set includes this instance. +// +// A caller answers "do I own this series?" with +// positions[SearchToken(tokens, key)], which is a binary search plus an array +// index. Both slices are immutable once returned, so the caller may hold on to +// them while classifying many series. +// +// The returned fingerprint changes if and only if ownership may have changed, so +// callers can use it to decide whether to redo work derived from the bitmap +// without having to compare the bitmap itself. It is zero before the ring has +// been read. +// +// Returns nil, nil before the ring has been read for the first time. Callers must +// treat that as "ownership unknown" rather than "owns nothing", because an +// ingester which has not yet seen the ring still holds its series. +func (i *Lifecycler) GetOwnedTokenPositions() (tokens []uint32, positions []bool, fingerprint uint64) { + i.countersLock.RLock() + defer i.countersLock.RUnlock() + return i.ownedTokens, i.ownedPositions, i.ownedFingerprint +} + func (i *Lifecycler) setTokens(tokens Tokens) { i.lifecyclerMetrics.tokensOwned.Set(float64(len(tokens))) @@ -938,7 +972,7 @@ func (i *Lifecycler) verifyTokens(ctx context.Context) bool { func (i *Lifecycler) compareTokens(fromRing Tokens) bool { sort.Sort(fromRing) - tokens := i.getTokens() + tokens := i.GetTokens() sort.Sort(tokens) if len(tokens) != len(fromRing) { @@ -970,14 +1004,14 @@ func (i *Lifecycler) autoJoin(ctx context.Context, targetState InstanceState, al // Need to make sure we didn't change the num of tokens configured myTokens, _ := ringDesc.TokensFor(i.ID) if !alreadyInRing { - myTokens = i.getTokens() + myTokens = i.GetTokens() } needTokens := i.cfg.NumTokens - len(myTokens) - if needTokens == 0 && myTokens.Equals(i.getTokens()) { + if needTokens == 0 && myTokens.Equals(i.GetTokens()) { // Tokens have been verified. No need to change them. state := i.GetState() - ringDesc.AddIngester(i.ID, i.Addr, i.Zone, i.getTokens(), state, i.getRegisteredAt()) + ringDesc.AddIngester(i.ID, i.Addr, i.Zone, i.GetTokens(), state, i.getRegisteredAt()) level.Info(i.logger).Log("msg", "auto joined with existing tokens", "ring", i.RingName, "state", state) return ringDesc, true, nil } @@ -993,7 +1027,7 @@ func (i *Lifecycler) autoJoin(ctx context.Context, targetState InstanceState, al i.setTokens(myTokens) state := i.GetState() - ringDesc.AddIngester(i.ID, i.Addr, i.Zone, i.getTokens(), state, i.getRegisteredAt()) + ringDesc.AddIngester(i.ID, i.Addr, i.Zone, i.GetTokens(), state, i.getRegisteredAt()) level.Info(i.logger).Log("msg", "auto joined with new tokens", "ring", i.RingName, "state", state) return ringDesc, true, nil @@ -1023,7 +1057,7 @@ func (i *Lifecycler) updateConsul(ctx context.Context) error { if !ok { // consul must have restarted level.Info(i.logger).Log("msg", "found empty ring, inserting tokens", "ring", i.RingName) - ringDesc.AddIngester(i.ID, i.Addr, i.Zone, i.getTokens(), i.GetState(), i.getRegisteredAt()) + ringDesc.AddIngester(i.ID, i.Addr, i.Zone, i.GetTokens(), i.GetState(), i.getRegisteredAt()) } else { instanceDesc.Timestamp = time.Now().Unix() instanceDesc.State = i.GetState() @@ -1092,6 +1126,48 @@ func (i *Lifecycler) updateCounters(ringDesc *Desc) { } } + // Recompute which ring token positions this instance owns, but only when the + // ring topology has actually changed. The walk is proportional to the total + // token count times the replication factor, which is far too expensive to + // repeat on every heartbeat, and heartbeats are by far the most common reason + // this function runs. + // + // This deliberately happens outside countersLock: HealthyInstancesCount is on + // the per-series ingestion path, and holding the lock across the walk would + // stall ingestion. + var ( + ownedFingerprint uint64 + ownedTokens []uint32 + ownedPositions []bool + ownedRecomputed bool + ) + + if ringDesc != nil { + ownedFingerprint = ownershipFingerprint(ringDesc) + + i.countersLock.RLock() + ownedUnchanged := i.ownedFingerprintOK && i.ownedFingerprint == ownedFingerprint + i.countersLock.RUnlock() + + if !ownedUnchanged { + var err error + ownedTokens, ownedPositions, err = OwnedTokenPositions( + ringDesc, + i.ID, + Write, + i.cfg.RingConfig.ReplicationFactor, + i.cfg.RingConfig.ZoneAwarenessEnabled, + ) + if err != nil { + // Keep the previous ownership rather than publishing an empty + // bitmap, which would read as "owns nothing" and under-count. + level.Error(i.logger).Log("msg", "failed to compute owned ring token positions", "ring", i.RingName, "err", err) + } else { + ownedRecomputed = true + } + } + } + zones := make([]string, 0, len(zonesMap)) for z := range zonesMap { zones = append(zones, z) @@ -1104,6 +1180,12 @@ func (i *Lifecycler) updateCounters(ringDesc *Desc) { i.healthyInstancesCount = healthyInstancesCount i.zonesCount = len(zones) i.zones = zones + if ownedRecomputed { + i.ownedTokens = ownedTokens + i.ownedPositions = ownedPositions + i.ownedFingerprint = ownedFingerprint + i.ownedFingerprintOK = true + } i.countersLock.Unlock() } diff --git a/pkg/ring/lifecycler_test.go b/pkg/ring/lifecycler_test.go index bc508370360..586f8b8f8a9 100644 --- a/pkg/ring/lifecycler_test.go +++ b/pkg/ring/lifecycler_test.go @@ -109,11 +109,11 @@ func TestLifecycler_RenewTokens(t *testing.T) { return nil }) - originalTokens := l1.getTokens() + originalTokens := l1.GetTokens() require.Len(t, originalTokens, 512) require.IsIncreasing(t, originalTokens) l1.RenewTokens(0.1, ctx) - newTokens := l1.getTokens() + newTokens := l1.GetTokens() require.Len(t, newTokens, 512) require.IsIncreasing(t, newTokens) diff := 0 @@ -364,7 +364,7 @@ func TestLifecycler_ShouldHandleInstanceAbruptlyRestarted(t *testing.T) { return checkNormalised(d, "ing1") }) - expectedTokens := l1.getTokens() + expectedTokens := l1.GetTokens() expectedRegisteredAt := l1.getRegisteredAt() // Wait 1 second because the registered timestamp has second precision. Without waiting @@ -382,7 +382,7 @@ func TestLifecycler_ShouldHandleInstanceAbruptlyRestarted(t *testing.T) { require.NoError(t, err) return checkNormalised(d, "ing1") && - expectedTokens.Equals(l2.getTokens()) && + expectedTokens.Equals(l2.GetTokens()) && expectedRegisteredAt.Unix() == l2.getRegisteredAt().Unix() }) } diff --git a/pkg/ring/owned_tokens.go b/pkg/ring/owned_tokens.go new file mode 100644 index 00000000000..4d632db7a64 --- /dev/null +++ b/pkg/ring/owned_tokens.go @@ -0,0 +1,118 @@ +package ring + +import ( + "encoding/binary" + "hash/fnv" + "sort" +) + +// ownershipFingerprint returns a value which changes if and only if something the +// replica-set walk depends on has changed: the set of instances, their zones, +// their states and their tokens. +// +// It deliberately ignores Timestamp and RegisteredTimestamp. Heartbeats bump +// Timestamp constantly, and recomputing ownership on every heartbeat would cost +// far more than making the per-series check cheap saves. +// +// Desc.RingCompare cannot be used for this: it reports a timestamp-only change +// and a state change as the same EqualButStatesAndTimestamps result, and +// ownership does depend on state, because a non-ACTIVE instance extends the +// replica set. +func ownershipFingerprint(d *Desc) uint64 { + ids := make([]string, 0, len(d.Ingesters)) + for id := range d.Ingesters { + ids = append(ids, id) + } + sort.Strings(ids) + + var ( + h = fnv.New64a() + buf [8]byte + sep = []byte{0} + ) + + for _, id := range ids { + instance := d.Ingesters[id] + + _, _ = h.Write([]byte(id)) + _, _ = h.Write(sep) + _, _ = h.Write([]byte(instance.Zone)) + _, _ = h.Write(sep) + + binary.LittleEndian.PutUint64(buf[:], uint64(instance.State)) + _, _ = h.Write(buf[:]) + + for _, token := range instance.Tokens { + binary.LittleEndian.PutUint32(buf[:4], token) + _, _ = h.Write(buf[:4]) + } + _, _ = h.Write(sep) + } + + return h.Sum64() +} + +// OwnedTokenPositions returns the ring's full sorted token list together with a +// parallel bitmap marking which token positions are owned by instanceID. +// +// A position p is "owned" by an instance when that instance is a member of the +// replica set the ring selects for any key k with SearchToken(tokens, k) == p. +// In other words ownership answers "would a write for a series hashing into this +// token range be replicated to me?", which is exactly the question the +// distributor answers when it routes. It is NOT the narrower question "am I the +// first instance in this token range": with a replication factor greater than +// one every series is held by several instances, and all of them must count it, +// because the per-instance series limit is itself scaled by the replication +// factor. +// +// The bitmap is derived from replicaSetAt, the same walk Ring.Get uses, so +// ownership and routing cannot drift apart. +// +// Ownership deliberately does NOT apply the replication strategy's health +// filter. It must be stable across heartbeat flapping: a series does not stop +// being this instance's responsibility because a peer missed a heartbeat. +// +// This is O(len(tokens) x replicationFactor) and is meant to be called once per +// ring change, never per series. Per-series questions are then answered with +// owned[SearchToken(tokens, key)], which is a binary search plus an array index. +func OwnedTokenPositions(d *Desc, instanceID string, op Operation, replicationFactor int, zoneAwarenessEnabled bool) ([]uint32, []bool, error) { + tokens := d.GetTokens() + owned := make([]bool, len(tokens)) + + // The ingester asks this question during startup, before it has read the + // ring, and briefly after being forgotten from the ring. Both are empty + // answers rather than errors. + if len(tokens) == 0 { + return tokens, owned, nil + } + + topology := ringTopology{ + tokens: tokens, + instanceByToken: d.getTokensInfo(), + instances: d.Ingesters, + numZones: len(d.getTokensByZone()), + replicationFactor: replicationFactor, + zoneAwarenessEnabled: zoneAwarenessEnabled, + } + + bufDescs, bufHosts, bufZones := MakeBuffersForGet() + + for position := range tokens { + _, instanceIDs, err := topology.replicaSetAt(position, op, bufDescs, bufHosts, bufZones) + if err != nil { + return nil, nil, err + } + + // Compare instance IDs, not addresses: InstanceDesc carries no ID of its + // own, and a ring is keyed by instance ID while Addr is a separate, + // different value. + for _, id := range instanceIDs { + if id == instanceID { + owned[position] = true + break + } + } + } + + return tokens, owned, nil +} diff --git a/pkg/ring/owned_tokens_test.go b/pkg/ring/owned_tokens_test.go new file mode 100644 index 00000000000..f061d7127ff --- /dev/null +++ b/pkg/ring/owned_tokens_test.go @@ -0,0 +1,438 @@ +package ring + +import ( + "fmt" + "math" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// ownershipTopologies is the matrix of ring shapes that ownership must be +// correct for. The zone-aware RF==numZones case is the one that happens to work +// with a "first token in my zone wins" ownership check, because there the ring +// picks exactly one instance per zone. Every other shape in this table does not, +// which is the point of the table. +var ownershipTopologies = map[string]struct { + numInstances int + numZones int + replicationFactor int + zoneAwarenessEnabled bool +}{ + // Zone awareness disabled. This is the Cortex default + // (-distributor.zone-awareness-enabled defaults to false) and the shape the + // ownership check must not get wrong. + "zones disabled, RF=1": {numInstances: 9, numZones: 3, replicationFactor: 1, zoneAwarenessEnabled: false}, + "zones disabled, RF=3": {numInstances: 9, numZones: 3, replicationFactor: 3, zoneAwarenessEnabled: false}, + "zones disabled, RF=5": {numInstances: 9, numZones: 3, replicationFactor: 5, zoneAwarenessEnabled: false}, + "single zone, RF=3": {numInstances: 9, numZones: 1, replicationFactor: 3, zoneAwarenessEnabled: false}, + "zones disabled, RF=12": {numInstances: 12, numZones: 3, replicationFactor: 12, zoneAwarenessEnabled: false}, + + // Zone awareness enabled. + "zone aware, RF == numZones": {numInstances: 9, numZones: 3, replicationFactor: 3, zoneAwarenessEnabled: true}, + "zone aware, RF > numZones": {numInstances: 9, numZones: 3, replicationFactor: 6, zoneAwarenessEnabled: true}, + "zone aware, RF < numZones": {numInstances: 9, numZones: 3, replicationFactor: 2, zoneAwarenessEnabled: true}, + "zone aware, RF % numZones != 0": { + numInstances: 12, numZones: 3, replicationFactor: 4, zoneAwarenessEnabled: true, + }, + "zone aware, 2 zones RF=3": {numInstances: 8, numZones: 2, replicationFactor: 3, zoneAwarenessEnabled: true}, +} + +// newRingForOwnershipTest builds a Ring in the same way the rest of this +// package's tests do, with every instance ACTIVE and heartbeating so that the +// replication strategy's Filter is an order-preserving identity. That makes +// Ring.Get observable as the raw replica-set walk, which is what ownership is +// defined against. +func newRingForOwnershipTest(ringDesc *Desc, replicationFactor int, zoneAwarenessEnabled bool) *Ring { + return &Ring{ + cfg: Config{ + HeartbeatTimeout: time.Hour, + ZoneAwarenessEnabled: zoneAwarenessEnabled, + ReplicationFactor: replicationFactor, + }, + ringDesc: ringDesc, + ringTokens: ringDesc.GetTokens(), + ringTokensByZone: ringDesc.getTokensByZone(), + ringInstanceByToken: ringDesc.getTokensInfo(), + ringZones: getZones(ringDesc.getTokensByZone()), + strategy: NewDefaultReplicationStrategy(), + KVClient: &MockClient{}, + } +} + +// keyForPosition returns a key that SearchToken maps to exactly position p, so +// that a per-position bitmap can be compared against a per-key Ring.Get. +func keyForPosition(tokens []uint32, p int) uint32 { + if p == 0 { + if tokens[0] > 0 { + // No token is greater than tokens[0]-1 before index 0. + return tokens[0] - 1 + } + // tokens[0] == 0, so instead rely on SearchToken wrapping past the end. + return math.MaxUint32 + } + // tokens[p-1] is not greater than itself, but tokens[p] is, so the first + // token strictly greater than this key sits at index p. + return tokens[p-1] +} + +// TestOwnedTokenPositions_MatchesGet is the load-bearing property for the whole +// owned-series feature: an instance owns a token position if and only if the +// ring would route a key in that position to it. +// +// If this test cannot be made to pass, then "the replica set is a pure function +// of the start index" is false and the owned-series design needs to change +// before anything is built on top of it. +func TestOwnedTokenPositions_MatchesGet(t *testing.T) { + for testName, testData := range ownershipTopologies { + t.Run(testName, func(t *testing.T) { + ringDesc := &Desc{Ingesters: generateRingInstances(testData.numInstances, testData.numZones, 128)} + ring := newRingForOwnershipTest(ringDesc, testData.replicationFactor, testData.zoneAwarenessEnabled) + + allTokens := ringDesc.GetTokens() + require.NotEmpty(t, allTokens) + + // Compute each instance's ownership bitmap once, the way a real + // caller would: once per ring change, not once per series. + bitmaps := make(map[string][]bool, len(ringDesc.Ingesters)) + for instanceID := range ringDesc.Ingesters { + tokens, owned, err := OwnedTokenPositions(ringDesc, instanceID, Write, testData.replicationFactor, testData.zoneAwarenessEnabled) + require.NoError(t, err) + require.Equal(t, allTokens, tokens, "returned token list must be the ring's full sorted token list") + require.Len(t, owned, len(allTokens), "bitmap must be parallel to the token list") + bitmaps[instanceID] = owned + } + + bufDescs, bufHosts, bufZones := MakeBuffersForGet() + + for p := range allTokens { + key := keyForPosition(allTokens, p) + require.Equal(t, p, SearchToken(allTokens, key), "test helper produced a key for the wrong position") + + set, err := ring.Get(key, Write, bufDescs, bufHosts, bufZones) + require.NoError(t, err) + + for instanceID, instanceDesc := range ringDesc.Ingesters { + // Compare via Addr because that is what ReplicationSet + // matches on, but drive the lookup from the instance ID, so + // this stays correct even when the two differ. + assert.Equal(t, set.Includes(instanceDesc.Addr), bitmaps[instanceID][p], + "ownership disagrees with routing at token position %d (token %d) for instance %s", + p, allTokens[p], instanceID) + } + } + }) + } +} + +// TestOwnedTokenPositions_Conservation checks the aggregate invariant that makes +// the series accounting add up: summed across all instances, every token +// position is owned exactly replicationFactor times. This is what guarantees +// that the per-ingester owned counts sum to RF x the total series count, and so +// that comparing owned against a limit already scaled by RF is dimensionally +// correct. +func TestOwnedTokenPositions_Conservation(t *testing.T) { + for testName, testData := range ownershipTopologies { + t.Run(testName, func(t *testing.T) { + ringDesc := &Desc{Ingesters: generateRingInstances(testData.numInstances, testData.numZones, 128)} + allTokens := ringDesc.GetTokens() + + ownersPerPosition := make([]int, len(allTokens)) + for instanceID := range ringDesc.Ingesters { + _, owned, err := OwnedTokenPositions(ringDesc, instanceID, Write, testData.replicationFactor, testData.zoneAwarenessEnabled) + require.NoError(t, err) + + for p, isOwned := range owned { + if isOwned { + ownersPerPosition[p]++ + } + } + } + + for p, owners := range ownersPerPosition { + assert.Equal(t, testData.replicationFactor, owners, + "token position %d (token %d) is owned by %d instances, expected exactly RF=%d", + p, allTokens[p], owners, testData.replicationFactor) + } + }) + } +} + +// TestOwnedTokenPositions_EmptyRing documents the behaviour on a ring with no +// instances. Callers must get an empty bitmap rather than a panic, because the +// ingester asks this question during startup before the ring has been read. +func TestOwnedTokenPositions_EmptyRing(t *testing.T) { + tokens, owned, err := OwnedTokenPositions(NewDesc(), "instance-1", Write, 3, false) + require.NoError(t, err) + assert.Empty(t, tokens) + assert.Empty(t, owned) +} + +// TestOwnedTokenPositions_AddrDiffersFromInstanceID guards a bug that the shared +// test fixtures cannot catch on their own: generateRingInstance happens to set +// Addr equal to the instance ID, so an implementation which compared addresses +// instead of instance IDs would pass every other test in this file and then own +// nothing at all in production, where a ring is keyed by instance ID +// ("ingester-zone-a-0") while Addr is a host:port. +func TestOwnedTokenPositions_AddrDiffersFromInstanceID(t *testing.T) { + const ( + numInstances = 9 + numTokens = 128 + replicationFactor = 3 + ) + + // Build a ring where the instance ID and the Addr are deliberately unrelated. + ringDesc := &Desc{Ingesters: map[string]InstanceDesc{}} + g := NewRandomTokenGenerator() + for i := 1; i <= numInstances; i++ { + instanceID := fmt.Sprintf("ingester-zone-a-%d", i) + addr := fmt.Sprintf("10.0.0.%d:9095", i) + ringDesc.Ingesters[instanceID] = InstanceDesc{ + Addr: addr, + Timestamp: time.Now().Unix(), + RegisteredTimestamp: time.Now().Unix(), + State: ACTIVE, + Tokens: g.GenerateTokens(NewDesc(), instanceID, "", numTokens, true), + Zone: "", + } + } + + ring := newRingForOwnershipTest(ringDesc, replicationFactor, false) + allTokens := ringDesc.GetTokens() + + bufDescs, bufHosts, bufZones := MakeBuffersForGet() + totalOwned := 0 + + for instanceID, instanceDesc := range ringDesc.Ingesters { + tokens, owned, err := OwnedTokenPositions(ringDesc, instanceID, Write, replicationFactor, false) + require.NoError(t, err) + require.Equal(t, allTokens, tokens) + + for p := range allTokens { + key := keyForPosition(allTokens, p) + require.Equal(t, p, SearchToken(allTokens, key)) + + set, err := ring.Get(key, Write, bufDescs, bufHosts, bufZones) + require.NoError(t, err) + + assert.Equal(t, set.Includes(instanceDesc.Addr), owned[p], + "ownership disagrees with routing at position %d for instance %s (addr %s)", + p, instanceID, instanceDesc.Addr) + + if owned[p] { + totalOwned++ + } + } + } + + // Sanity check that the test actually exercised ownership rather than + // trivially agreeing on an all-false bitmap. + assert.Equal(t, replicationFactor*len(allTokens), totalOwned) +} + +// TestOwnedTokenPositions_NonActiveInstancesExtendReplicaSet pins down what +// ownership means while instances are transitioning, which is deliberately not the +// same as what the write path routes. +// +// An instance in a state other than ACTIVE extends the replica set, so the ring +// picks an additional instance for the affected token ranges. Both the +// transitioning instance and the extra one own those ranges, because both really +// are holding that data: the departing instance still has the series in its head, +// and the distributor really is writing to the extension. +// +// Ownership is computed from the raw walk and deliberately does not apply the +// replication strategy's health filter, so a LEAVING instance keeps owning its +// series even though Ring.Get for a write excludes it. That asymmetry is the point: +// ownership answers "what am I holding, and must therefore count against my +// limit", not "where would a new write go". Filtering here would make an +// instance's own series count depend on its peers' heartbeat luck. +func TestOwnedTokenPositions_NonActiveInstancesExtendReplicaSet(t *testing.T) { + const ( + numInstances = 9 + replicationFactor = 3 + transitioning = "instance-1" + ) + + for _, state := range []InstanceState{JOINING, LEAVING, READONLY} { + t.Run(state.String(), func(t *testing.T) { + ringDesc := &Desc{Ingesters: generateRingInstances(numInstances, 1, 128)} + + instance := ringDesc.Ingesters[transitioning] + instance.State = state + ringDesc.Ingesters[transitioning] = instance + + ring := newRingForOwnershipTest(ringDesc, replicationFactor, false) + allTokens := ringDesc.GetTokens() + + ownersPerPosition := make([]int, len(allTokens)) + var transitioningOwned []bool + + for instanceID := range ringDesc.Ingesters { + _, owned, err := OwnedTokenPositions(ringDesc, instanceID, Write, replicationFactor, false) + require.NoError(t, err) + + if instanceID == transitioning { + transitioningOwned = owned + } + for p, isOwned := range owned { + if isOwned { + ownersPerPosition[p]++ + } + } + } + + // The transitioning instance still holds its series, so it still owns them. + transitioningOwnedCount := 0 + for _, isOwned := range transitioningOwned { + if isOwned { + transitioningOwnedCount++ + } + } + assert.NotZero(t, transitioningOwnedCount, + "a %s instance still holds its series and must still own them", state) + + bufDescs, bufHosts, bufZones := MakeBuffersForGet() + extended := 0 + + for p := range allTokens { + // Extension means these ranges have more owners than the replication + // factor, never fewer. + assert.GreaterOrEqual(t, ownersPerPosition[p], replicationFactor, + "token position %d has fewer owners than RF", p) + if ownersPerPosition[p] > replicationFactor { + extended++ + } + + key := keyForPosition(allTokens, p) + set, err := ring.Get(key, Write, bufDescs, bufHosts, bufZones) + require.NoError(t, err) + + // Writes still land on exactly RF healthy instances, and never on the + // transitioning one, even where it owns the range. + assert.Len(t, set.Instances, replicationFactor) + assert.False(t, set.Includes(ringDesc.Ingesters[transitioning].Addr), + "a %s instance must not receive writes at position %d", state, p) + } + + assert.NotZero(t, extended, + "a %s instance must cause the replica set to be extended somewhere", state) + }) + } +} + +// TestOwnershipFingerprint pins down the fingerprint's two jobs. It must change +// whenever ownership could have changed, because a false "unchanged" leaves every +// ingester using a stale ownership bitmap. It must NOT change on a heartbeat, +// because heartbeats are the overwhelmingly common reason the lifecycler reruns +// this, and recomputing ownership each time would cost far more than the cheap +// per-series lookup saves. +func TestOwnershipFingerprint(t *testing.T) { + base := &Desc{Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: ACTIVE, Tokens: []uint32{100, 300}, Timestamp: 1000}, + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200, 400}, Timestamp: 1000}, + }} + baseFingerprint := ownershipFingerprint(base) + + t.Run("is stable for equal rings", func(t *testing.T) { + same := &Desc{Ingesters: map[string]InstanceDesc{ + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200, 400}, Timestamp: 1000}, + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: ACTIVE, Tokens: []uint32{100, 300}, Timestamp: 1000}, + }} + assert.Equal(t, baseFingerprint, ownershipFingerprint(same), + "fingerprint must not depend on map iteration order") + }) + + t.Run("ignores heartbeats", func(t *testing.T) { + heartbeat := &Desc{Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: ACTIVE, Tokens: []uint32{100, 300}, Timestamp: 9999}, + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200, 400}, Timestamp: 8888}, + }} + assert.Equal(t, baseFingerprint, ownershipFingerprint(heartbeat), + "a heartbeat must not force an ownership recompute") + }) + + // Each of these changes the replica set, so each must invalidate the bitmap. + // The state case is the one Desc.RingCompare cannot distinguish from a + // heartbeat, which is why this fingerprint exists at all. + for name, changed := range map[string]*Desc{ + "state change": {Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: LEAVING, Tokens: []uint32{100, 300}, Timestamp: 1000}, + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200, 400}, Timestamp: 1000}, + }}, + "token change": {Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: ACTIVE, Tokens: []uint32{100, 301}, Timestamp: 1000}, + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200, 400}, Timestamp: 1000}, + }}, + "zone change": {Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-c", State: ACTIVE, Tokens: []uint32{100, 300}, Timestamp: 1000}, + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200, 400}, Timestamp: 1000}, + }}, + "instance added": {Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: ACTIVE, Tokens: []uint32{100, 300}, Timestamp: 1000}, + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200, 400}, Timestamp: 1000}, + "instance-3": {Addr: "10.0.0.3:9095", Zone: "zone-c", State: ACTIVE, Tokens: []uint32{500}, Timestamp: 1000}, + }}, + "instance removed": {Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: ACTIVE, Tokens: []uint32{100, 300}, Timestamp: 1000}, + }}, + "tokens moved between instances": {Ingesters: map[string]InstanceDesc{ + "instance-1": {Addr: "10.0.0.1:9095", Zone: "zone-a", State: ACTIVE, Tokens: []uint32{100, 300, 400}, Timestamp: 1000}, + "instance-2": {Addr: "10.0.0.2:9095", Zone: "zone-b", State: ACTIVE, Tokens: []uint32{200}, Timestamp: 1000}, + }}, + } { + t.Run("changes on "+name, func(t *testing.T) { + assert.NotEqual(t, baseFingerprint, ownershipFingerprint(changed)) + }) + } +} + +// BenchmarkOwnedTokenPositions measures the cost of the once-per-ring-change +// precompute. This is the price paid to make the per-series ownership check a +// binary search plus an array index, so it is expected to be large relative to a +// single Ring.Get and still negligible relative to the interval between ring +// changes. +func BenchmarkOwnedTokenPositions(b *testing.B) { + const ( + numInstances = 100 + numZones = 3 + replicationFactor = 3 + ) + + ringDesc := &Desc{Ingesters: generateRingInstances(numInstances, numZones, numTokens)} + + var instanceID string + for id := range ringDesc.Ingesters { + instanceID = id + break + } + + b.ReportAllocs() + + for b.Loop() { + _, owned, err := OwnedTokenPositions(ringDesc, instanceID, Write, replicationFactor, true) + if err != nil { + b.Fatal(err) + } + if len(owned) == 0 { + b.Fatal("expected a non-empty bitmap") + } + } +} + +// TestOwnedTokenPositions_UnknownInstance documents that an instance which is +// not in the ring owns nothing. The ingester can be in this state briefly after +// being forgotten from the ring, and it must not then claim ownership of +// everything. +func TestOwnedTokenPositions_UnknownInstance(t *testing.T) { + ringDesc := &Desc{Ingesters: generateRingInstances(6, 3, 128)} + + tokens, owned, err := OwnedTokenPositions(ringDesc, "instance-does-not-exist", Write, 3, false) + require.NoError(t, err) + require.Len(t, owned, len(tokens)) + + for p, isOwned := range owned { + assert.False(t, isOwned, "unknown instance must not own token position %d", p) + } +} diff --git a/pkg/ring/ring.go b/pkg/ring/ring.go index b510a907625..f2419a80ebe 100644 --- a/pkg/ring/ring.go +++ b/pkg/ring/ring.go @@ -392,24 +392,70 @@ func (r *Ring) updateRingState(ringDesc *Desc) { r.updateRingMetrics(rc) } -// Get returns n (or more) instances which form the replicas for the given key. -// This implementation guarantees: -// - Stability: given the same ring, two invocations returns the same set for same operation. -// - Consistency: adding/removing 1 instance from the ring returns set with no more than 1 difference for same operation. -func (r *Ring) Get(key uint32, op Operation, bufDescs []InstanceDesc, bufHosts []string, bufZones map[string]int) (ReplicationSet, error) { - r.mtx.RLock() - defer r.mtx.RUnlock() - if r.ringDesc == nil || len(r.ringTokens) == 0 { - return ReplicationSet{}, ErrEmptyRing +// ringTopology is the read-only view of a ring required to compute a replica +// set. It exists so that the replica-set walk can be shared between Ring.Get, +// which starts at the token position of a single key, and ownership +// computation, which walks every token position in turn. Sharing one walk is +// what guarantees that an instance's idea of which series it owns cannot drift +// from the ring's idea of where those series are routed. +type ringTopology struct { + // tokens is the full, sorted token list of the ring. + tokens []uint32 + // instanceByToken maps each token to information about the instance holding it. + instanceByToken map[uint32]instanceInfo + // instances holds the instance descriptors, keyed by instance ID. + instances map[string]InstanceDesc + // numZones is the number of distinct zones holding at least one instance. + numZones int + + replicationFactor int + zoneAwarenessEnabled bool +} + +// topology returns the read-only view of the ring used by the replica-set walk. +// Callers must hold at least a read lock on r.mtx. This is cheap: it copies +// slice and map headers, never their contents. +func (r *Ring) topology() ringTopology { + return ringTopology{ + tokens: r.ringTokens, + instanceByToken: r.ringInstanceByToken, + instances: r.ringDesc.Ingesters, + numZones: len(r.ringZones), + replicationFactor: r.cfg.ReplicationFactor, + zoneAwarenessEnabled: r.cfg.ZoneAwarenessEnabled, + } +} + +// replicaSetAt walks the ring starting from the given token position and returns +// the instances which form the replica set for any key mapping to that +// position. It is the single definition of "which instances hold this data". +// +// Two parallel slices are returned: the instance descriptors, and the instance +// IDs which key them in the ring. Both are in walk order. The IDs are returned +// because InstanceDesc does not carry its own ID, and an instance's ID is not +// interchangeable with its Addr, so a caller asking "am I in this set?" can only +// answer correctly by comparing IDs. +// +// The returned instances have NOT been health filtered. Callers which need only +// reachable instances apply a ReplicationStrategy afterwards; callers which need +// an answer that stays stable across heartbeat flapping use the result directly. +// +// Both returned slices alias the caller's buffers and are only valid until the +// next call sharing those buffers. +func (t ringTopology) replicaSetAt(start int, op Operation, bufDescs []InstanceDesc, bufHosts []string, bufZones map[string]int) ([]InstanceDesc, []string, error) { + if len(t.tokens) == 0 { + return nil, nil, ErrEmptyRing } var ( - replicationFactor = r.cfg.ReplicationFactor - instances = bufDescs[:0] - start = searchToken(r.ringTokens, key) - iterations = 0 - maxInstancePerZone = replicationFactor / len(r.ringZones) - zonesWithExtraInstance = replicationFactor % len(r.ringZones) + replicationFactor = t.replicationFactor + instances = bufDescs[:0] + iterations = 0 + + // With no zones known there is no per-zone cap to apply, so allow the + // whole replica set to come from the single (unnamed) zone. + maxInstancePerZone = replicationFactor + zonesWithExtraInstance = 0 // We use a slice instead of a map because it's faster to search within a // slice than lookup a map for a very low number of items. @@ -417,16 +463,21 @@ func (r *Ring) Get(key uint32, op Operation, bufDescs []InstanceDesc, bufHosts [ numOfInstanceByZone = resetZoneMap(bufZones) ) - for i := start; len(distinctHosts) < replicationFactor && iterations < len(r.ringTokens); i++ { + if t.numZones > 0 { + maxInstancePerZone = replicationFactor / t.numZones + zonesWithExtraInstance = replicationFactor % t.numZones + } + + for i := start; len(distinctHosts) < replicationFactor && iterations < len(t.tokens); i++ { iterations++ // Wrap i around in the ring. - i %= len(r.ringTokens) - token := r.ringTokens[i] + i %= len(t.tokens) + token := t.tokens[i] - info, ok := r.ringInstanceByToken[token] + info, ok := t.instanceByToken[token] if !ok { // This should never happen unless a bug in the ring code. - return ReplicationSet{}, ErrInconsistentTokensInfo + return nil, nil, ErrInconsistentTokensInfo } // We want n *distinct* instances. @@ -435,7 +486,7 @@ func (r *Ring) Get(key uint32, op Operation, bufDescs []InstanceDesc, bufHosts [ } // Ignore if the instances don't have a zone set. - if r.cfg.ZoneAwarenessEnabled && info.Zone != "" { + if t.zoneAwarenessEnabled && info.Zone != "" { maxNumOfInstance := maxInstancePerZone // If we still have room for zones with extra instance, increase the instance threshold by 1 if zonesWithExtraInstance > 0 { @@ -448,13 +499,13 @@ func (r *Ring) Get(key uint32, op Operation, bufDescs []InstanceDesc, bufHosts [ } distinctHosts = append(distinctHosts, info.InstanceID) - instance := r.ringDesc.Ingesters[info.InstanceID] + instance := t.instances[info.InstanceID] // Check whether the replica set should be extended given we're including // this instance. if op.ShouldExtendReplicaSetOnState(instance.State) { replicationFactor++ - } else if r.cfg.ZoneAwarenessEnabled && info.Zone != "" { + } else if t.zoneAwarenessEnabled && info.Zone != "" { // We should only add the zone if we are not going to extend, // as we want to extend the instance in the same AZ. if numOfInstance, ok := numOfInstanceByZone[info.Zone]; !ok { @@ -471,6 +522,25 @@ func (r *Ring) Get(key uint32, op Operation, bufDescs []InstanceDesc, bufHosts [ instances = append(instances, instance) } + return instances, distinctHosts, nil +} + +// Get returns n (or more) instances which form the replicas for the given key. +// This implementation guarantees: +// - Stability: given the same ring, two invocations returns the same set for same operation. +// - Consistency: adding/removing 1 instance from the ring returns set with no more than 1 difference for same operation. +func (r *Ring) Get(key uint32, op Operation, bufDescs []InstanceDesc, bufHosts []string, bufZones map[string]int) (ReplicationSet, error) { + r.mtx.RLock() + defer r.mtx.RUnlock() + if r.ringDesc == nil || len(r.ringTokens) == 0 { + return ReplicationSet{}, ErrEmptyRing + } + + instances, _, err := r.topology().replicaSetAt(SearchToken(r.ringTokens, key), op, bufDescs, bufHosts, bufZones) + if err != nil { + return ReplicationSet{}, err + } + healthyInstances, maxFailure, err := r.strategy.Filter(instances, op, r.cfg.ReplicationFactor, r.cfg.HeartbeatTimeout, r.cfg.ZoneAwarenessEnabled, r.KVClient.LastUpdateTime(r.key)) if err != nil { return ReplicationSet{}, err @@ -863,7 +933,7 @@ func (r *Ring) shuffleShard(identifier string, size int, lookbackPeriod time.Dur finalInstancesPerZone++ } for i := 0; i < finalInstancesPerZone; i++ { - start := searchToken(tokens, random.Uint32()) + start := SearchToken(tokens, random.Uint32()) iterations := 0 found := false diff --git a/pkg/ring/token.go b/pkg/ring/token.go new file mode 100644 index 00000000000..87166495d06 --- /dev/null +++ b/pkg/ring/token.go @@ -0,0 +1,56 @@ +package ring + +import ( + "github.com/cortexproject/cortex/pkg/cortexpb" + "github.com/cortexproject/cortex/pkg/util" + "github.com/cortexproject/cortex/pkg/util/extract" +) + +// TokenForLabels returns the ring token (hash key) for a set of series labels. +// This determines which ingester in the ring is responsible for this series. +// Used by both the distributor (to route) and the ingester (to check ownership). +func TokenForLabels(userID string, labels []cortexpb.LabelAdapter, shouldShardByAllLabels bool) (uint32, error) { + if shouldShardByAllLabels { + return ShardByAllLabels(userID, labels), nil + } + + unsafeMetricName, err := extract.UnsafeMetricNameFromLabelAdapters(labels) + if err != nil { + return 0, err + } + return ShardByMetricName(userID, unsafeMetricName), nil +} + +// TokenForMetadata returns the ring token for metadata routing. +func TokenForMetadata(userID string, metricName string, shouldShardByAllLabels bool) uint32 { + if shouldShardByAllLabels { + return ShardByMetricName(userID, metricName) + } + return shardByUser(userID) +} + +// ShardByAllLabels generates a token from userID + all label name/value pairs. +// This function generates different values for different order of same labels. +func ShardByAllLabels(userID string, labels []cortexpb.LabelAdapter) uint32 { + h := shardByUser(userID) + for _, label := range labels { + if len(label.Value) > 0 { + h = util.HashAdd32(h, label.Name) + h = util.HashAdd32(h, label.Value) + } + } + return h +} + +// ShardByMetricName returns the token for the given metric. +func ShardByMetricName(userID string, metricName string) uint32 { + h := shardByUser(userID) + h = util.HashAdd32(h, metricName) + return h +} + +func shardByUser(userID string) uint32 { + h := util.HashNew32() + h = util.HashAdd32(h, userID) + return h +} diff --git a/pkg/ring/util.go b/pkg/ring/util.go index 66a176c0543..608032e12e0 100644 --- a/pkg/ring/util.go +++ b/pkg/ring/util.go @@ -139,8 +139,8 @@ func getZones(tokens map[string][]uint32) []string { return zones } -// searchToken returns the offset of the tokens entry holding the range for the provided key. -func searchToken(tokens []uint32, key uint32) int { +// SearchToken returns the offset of the tokens entry holding the range for the provided key. +func SearchToken(tokens []uint32, key uint32) int { i := sort.Search(len(tokens), func(x int) bool { return tokens[x] > key }) diff --git a/pkg/storage/tsdb/util.go b/pkg/storage/tsdb/util.go index 8b10403c09e..ae9ec8c1ed7 100644 --- a/pkg/storage/tsdb/util.go +++ b/pkg/storage/tsdb/util.go @@ -3,15 +3,15 @@ package tsdb import ( "github.com/oklog/ulid/v2" - "github.com/cortexproject/cortex/pkg/ingester/client" + "github.com/cortexproject/cortex/pkg/util" ) // HashBlockID returns a 32-bit hash of the block ID useful for // ring-based sharding. func HashBlockID(id ulid.ULID) uint32 { - h := client.HashNew32() + h := util.HashNew32() for _, b := range id { - h = client.HashAddByte32(h, b) + h = util.HashAddByte32(h, b) } return h } diff --git a/pkg/ingester/client/fnv.go b/pkg/util/fnv.go similarity index 66% rename from pkg/ingester/client/fnv.go rename to pkg/util/fnv.go index fd35a174a1a..d5e25941ab1 100644 --- a/pkg/ingester/client/fnv.go +++ b/pkg/util/fnv.go @@ -12,9 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. -package client +package util -// Inline and byte-free variant of hash/fnv's fnv64a. +// Inline and byte-free variant of hash/fnv's fnv64a and fnv32. const ( offset64 = 14695981039346656037 @@ -23,14 +23,14 @@ const ( prime32 = 16777619 ) -// hashNew initializes a new fnv64a hash value. -func hashNew() uint64 { +// HashNew initializes a new fnv64a hash value. +func HashNew() uint64 { return offset64 } -// hashAdd adds a string to a fnv64a hash value, returning the updated hash. +// HashAdd adds a string to a fnv64a hash value, returning the updated hash. // Note this is the same algorithm as Go stdlib `sum64a.Write()` -func hashAdd(h uint64, s string) uint64 { +func HashAdd(h uint64, s string) uint64 { for i := 0; i < len(s); i++ { h ^= uint64(s[i]) h *= prime64 @@ -38,14 +38,21 @@ func hashAdd(h uint64, s string) uint64 { return h } -// hashAddByte adds a byte to a fnv64a hash value, returning the updated hash. -func hashAddByte(h uint64, b byte) uint64 { +// HashAddByte adds a byte to a fnv64a hash value, returning the updated hash. +func HashAddByte(h uint64, b byte) uint64 { h ^= uint64(b) h *= prime64 return h } -// HashNew32 initializies a new fnv32 hash value. +// HashAddUint adds a uint64 to a fnv64a hash value, returning the updated hash. +func HashAddUint(h uint64, i uint64) uint64 { + h ^= i + h *= prime64 + return h +} + +// HashNew32 initializes a new fnv32 hash value. func HashNew32() uint32 { return offset32 } @@ -66,3 +73,10 @@ func HashAddByte32(h uint32, b byte) uint32 { h ^= uint32(b) return h } + +// HashAddUint32 adds a uint32 to a fnv32 hash value, returning the updated hash. +func HashAddUint32(h uint32, i uint32) uint32 { + h *= prime32 + h ^= i + return h +} diff --git a/schemas/cortex-config-schema.json b/schemas/cortex-config-schema.json index 0178c592d20..fa7e6701ad4 100644 --- a/schemas/cortex-config-schema.json +++ b/schemas/cortex-config-schema.json @@ -5414,6 +5414,18 @@ "x-cli-flag": "ingester.metadata-retain-period", "x-format": "duration" }, + "owned_series_limit_enforcement_enabled": { + "default": false, + "description": "Use owned series count for limit enforcement. Requires owned-series-metrics-enabled. When enabled, PreCreation uses owned count instead of Head().NumSeries() for both per-user and instance-level limits.", + "type": "boolean", + "x-cli-flag": "ingester.owned-series-limit-enforcement-enabled" + }, + "owned_series_metrics_enabled": { + "default": false, + "description": "Enable tracking of owned series per user. When enabled, the ingester computes series ownership based on the ring and emits cortex_ingester_owned_series metric.", + "type": "boolean", + "x-cli-flag": "ingester.owned-series-metrics-enabled" + }, "query_protection": { "properties": { "eviction": {