-
Notifications
You must be signed in to change notification settings - Fork 460
/
blocks_store_replicated_set.go
157 lines (126 loc) · 4.74 KB
/
blocks_store_replicated_set.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
// SPDX-License-Identifier: AGPL-3.0-only
// Provenance-includes-location: https://github.com/cortexproject/cortex/blob/master/pkg/querier/blocks_store_replicated_set.go
// Provenance-includes-license: Apache-2.0
// Provenance-includes-copyright: The Cortex Authors.
package querier
import (
"context"
"fmt"
"math/rand"
"github.com/go-kit/log"
"github.com/grafana/dskit/ring"
"github.com/grafana/dskit/ring/client"
"github.com/grafana/dskit/services"
"github.com/oklog/ulid"
"github.com/pkg/errors"
"github.com/prometheus/client_golang/prometheus"
mimir_tsdb "github.com/grafana/mimir/pkg/storage/tsdb"
"github.com/grafana/mimir/pkg/storegateway"
"github.com/grafana/mimir/pkg/util"
)
type loadBalancingStrategy int
const (
noLoadBalancing = loadBalancingStrategy(iota)
randomLoadBalancing
)
// BlocksStoreSet implementation used when the blocks are sharded and replicated across
// a set of store-gateway instances.
type blocksStoreReplicationSet struct {
services.Service
storesRing *ring.Ring
clientsPool *client.Pool
balancingStrategy loadBalancingStrategy
limits BlocksStoreLimits
// Subservices manager.
subservices *services.Manager
subservicesWatcher *services.FailureWatcher
}
func newBlocksStoreReplicationSet(
storesRing *ring.Ring,
balancingStrategy loadBalancingStrategy,
limits BlocksStoreLimits,
clientConfig ClientConfig,
logger log.Logger,
reg prometheus.Registerer,
) (*blocksStoreReplicationSet, error) {
s := &blocksStoreReplicationSet{
storesRing: storesRing,
clientsPool: newStoreGatewayClientPool(client.NewRingServiceDiscovery(storesRing), clientConfig, logger, reg),
balancingStrategy: balancingStrategy,
limits: limits,
subservicesWatcher: services.NewFailureWatcher(),
}
var err error
s.subservices, err = services.NewManager(s.storesRing, s.clientsPool)
if err != nil {
return nil, err
}
s.Service = services.NewBasicService(s.starting, s.running, s.stopping)
return s, nil
}
func (s *blocksStoreReplicationSet) starting(ctx context.Context) error {
s.subservicesWatcher.WatchManager(s.subservices)
if err := services.StartManagerAndAwaitHealthy(ctx, s.subservices); err != nil {
return errors.Wrap(err, "unable to start blocks store set subservices")
}
return nil
}
func (s *blocksStoreReplicationSet) running(ctx context.Context) error {
for {
select {
case <-ctx.Done():
return nil
case err := <-s.subservicesWatcher.Chan():
return errors.Wrap(err, "blocks store set subservice failed")
}
}
}
func (s *blocksStoreReplicationSet) stopping(_ error) error {
return services.StopManagerAndAwaitStopped(context.Background(), s.subservices)
}
func (s *blocksStoreReplicationSet) GetClientsFor(userID string, blockIDs []ulid.ULID, exclude map[ulid.ULID][]string) (map[BlocksStoreClient][]ulid.ULID, error) {
blocks := make(map[string][]ulid.ULID)
instances := make(map[string]ring.InstanceDesc)
userRing := storegateway.GetShuffleShardingSubring(s.storesRing, userID, s.limits)
// Find the replication set of each block we need to query.
for _, blockID := range blockIDs {
// Do not reuse the same buffer across multiple Get() calls because we do retain the
// returned replication set.
bufDescs, bufHosts, bufZones := ring.MakeBuffersForGet()
set, err := userRing.Get(mimir_tsdb.HashBlockID(blockID), storegateway.BlocksRead, bufDescs, bufHosts, bufZones)
if err != nil {
return nil, errors.Wrapf(err, "failed to get store-gateway replication set owning the block %s", blockID.String())
}
// Pick a non excluded store-gateway instance.
inst := getNonExcludedInstance(set, exclude[blockID], s.balancingStrategy)
if inst == nil {
return nil, fmt.Errorf("no store-gateway instance left after checking exclude for block %s", blockID.String())
}
instances[inst.Addr] = *inst
blocks[inst.Addr] = append(blocks[inst.Addr], blockID)
}
clients := map[BlocksStoreClient][]ulid.ULID{}
// Get the client for each store-gateway.
for addr, instance := range instances {
c, err := s.clientsPool.GetClientForInstance(instance)
if err != nil {
return nil, errors.Wrapf(err, "failed to get store-gateway client for %s %s", instance.Id, addr)
}
clients[c.(BlocksStoreClient)] = blocks[addr]
}
return clients, nil
}
func getNonExcludedInstance(set ring.ReplicationSet, exclude []string, balancingStrategy loadBalancingStrategy) *ring.InstanceDesc {
if balancingStrategy == randomLoadBalancing {
// Randomize the list of instances to not always query the same one.
rand.Shuffle(len(set.Instances), func(i, j int) {
set.Instances[i], set.Instances[j] = set.Instances[j], set.Instances[i]
})
}
for _, instance := range set.Instances {
if !util.StringsContain(exclude, instance.Addr) {
return &instance
}
}
return nil
}