Skip to content
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions docs/configuration/config-file-reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -4198,6 +4198,18 @@ lifecycler:
# CLI flag: -ingester.head-queried-series-metrics-windows
[head_queried_series_metrics_windows: <list of duration> | 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: <boolean> | 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: <boolean> | default = false]

# Enable uploading compacted blocks.
# CLI flag: -ingester.upload-compacted-blocks-enabled
[upload_compacted_blocks_enabled: <boolean> | default = true]
Expand Down
4 changes: 4 additions & 0 deletions docs/configuration/v1-guarantees.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
43 changes: 2 additions & 41 deletions pkg/distributor/distributor.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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.
Expand Down
8 changes: 4 additions & 4 deletions pkg/distributor/distributor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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{}{}
Expand Down Expand Up @@ -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"},
Expand Down
2 changes: 1 addition & 1 deletion pkg/distributor/query.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}

Expand Down
Loading