unified-storage: expose ring replication factor config (#106345)
* config ring replication factor * change default * rename * fix test * fix
This commit is contained in:
+1
-1
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -576,6 +576,7 @@ type Cfg struct {
|
||||
MemberlistJoinMember string
|
||||
MemberlistClusterLabel string
|
||||
MemberlistClusterLabelVerificationDisabled bool
|
||||
SearchRingReplicationFactor int
|
||||
InstanceID string
|
||||
SprinklesApiServer string
|
||||
SprinklesApiServerPageLimit int
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user