diff --git a/pkg/server/ring.go b/pkg/server/ring.go index cd89f719f8a..90026e58391 100644 --- a/pkg/server/ring.go +++ b/pkg/server/ring.go @@ -91,7 +91,7 @@ func toRingConfig(cfg *setting.Cfg, KVStore kv.Config) ring.Config { rc.KVStore = KVStore rc.HeartbeatTimeout = resource.RingHeartbeatTimeout - rc.ReplicationFactor = 1 + rc.ReplicationFactor = cfg.SearchRingReplicationFactor return rc } diff --git a/pkg/server/search_server_distributor_test.go b/pkg/server/search_server_distributor_test.go index 32dddffcec8..1fce76d8936 100644 --- a/pkg/server/search_server_distributor_test.go +++ b/pkg/server/search_server_distributor_test.go @@ -274,6 +274,7 @@ func initDistributorServerForTest(t *testing.T, memberlistPort int) testModuleSe cfg.MemberlistJoinMember = "127.0.0.1:" + strconv.Itoa(memberlistPort) cfg.MemberlistAdvertiseAddr = "127.0.0.1" cfg.MemberlistAdvertisePort = memberlistPort + cfg.SearchRingReplicationFactor = 1 cfg.Target = []string{modules.SearchServerDistributor} cfg.InstanceID = "distributor" // does nothing for the distributor but may be useful to debug tests @@ -308,6 +309,7 @@ func createStorageServerApi(t *testing.T, instanceId int, dbType, dbConnStr stri cfg.MemberlistJoinMember = "127.0.0.1:" + strconv.Itoa(memberlistPort) cfg.MemberlistAdvertiseAddr = "127.0.0.1" cfg.MemberlistAdvertisePort = getRandomPort() + cfg.SearchRingReplicationFactor = 1 cfg.InstanceID = "instance-" + strconv.Itoa(instanceId) cfg.IndexPath = t.TempDir() + cfg.InstanceID cfg.IndexFileThreshold = testIndexFileThreshold diff --git a/pkg/setting/setting.go b/pkg/setting/setting.go index f14dbc03fc7..32f7be6830b 100644 --- a/pkg/setting/setting.go +++ b/pkg/setting/setting.go @@ -576,6 +576,7 @@ type Cfg struct { MemberlistJoinMember string MemberlistClusterLabel string MemberlistClusterLabelVerificationDisabled bool + SearchRingReplicationFactor int InstanceID string SprinklesApiServer string SprinklesApiServerPageLimit int diff --git a/pkg/setting/setting_unified_storage.go b/pkg/setting/setting_unified_storage.go index f0881cf95db..2e3a4a3122b 100644 --- a/pkg/setting/setting_unified_storage.go +++ b/pkg/setting/setting_unified_storage.go @@ -65,6 +65,7 @@ func (cfg *Cfg) setUnifiedStorageConfig() { cfg.MemberlistJoinMember = section.Key("memberlist_join_member").String() cfg.MemberlistClusterLabel = section.Key("memberlist_cluster_label").String() cfg.MemberlistClusterLabelVerificationDisabled = section.Key("memberlist_cluster_label_verification_disabled").MustBool(false) + cfg.SearchRingReplicationFactor = section.Key("search_ring_replication_factor").MustInt(1) cfg.InstanceID = section.Key("instance_id").String() cfg.IndexFileThreshold = section.Key("index_file_threshold").MustInt(10) cfg.IndexMinCount = section.Key("index_min_count").MustInt(1) diff --git a/pkg/storage/unified/resource/search_server_distributor.go b/pkg/storage/unified/resource/search_server_distributor.go index 9637e3055cb..c148448ca7b 100644 --- a/pkg/storage/unified/resource/search_server_distributor.go +++ b/pkg/storage/unified/resource/search_server_distributor.go @@ -81,7 +81,7 @@ type distributorServer struct { log log.Logger } -var ringOp = ring.NewOp([]ring.InstanceState{ring.ACTIVE}, func(s ring.InstanceState) bool { +var activeRingOp = ring.NewOp([]ring.InstanceState{ring.ACTIVE}, func(s ring.InstanceState) bool { return s != ring.ACTIVE }) @@ -128,7 +128,7 @@ func (ds *distributorServer) getClientToDistributeRequest(ctx context.Context, n return ctx, nil, err } - rs, err := ds.ring.Get(ringHasher.Sum32(), ringOp, nil, nil, nil) + rs, err := ds.ring.GetWithOptions(ringHasher.Sum32(), activeRingOp, ring.WithReplicationFactor(ds.ring.ReplicationFactor())) if err != nil { return ctx, nil, err }