-
Notifications
You must be signed in to change notification settings - Fork 2.1k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
query: add memcached autodiscovery support (#4487)
* add memcached autodiscovery support Signed-off-by: Roy Chiang <[email protected]> * pass logger to provider struct Signed-off-by: Roy Chiang <[email protected]> * fix linter issues Signed-off-by: Roy Chiang <[email protected]> * remove copypasted comment Signed-off-by: Roy Chiang <[email protected]> * address comments Signed-off-by: Roy Chiang <[email protected]> * Update docs/components/store.md Co-authored-by: Giedrius Statkevičius <[email protected]> Signed-off-by: Roy Chiang <[email protected]> * add lock/unlock for autodiscovery provider Signed-off-by: Roy Chiang <[email protected]> Co-authored-by: Giedrius Statkevičius <[email protected]>
- Loading branch information
1 parent
aee7c55
commit baea4ce
Showing
8 changed files
with
435 additions
and
17 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,114 @@ | ||
// Copyright (c) The Thanos Authors. | ||
// Licensed under the Apache License 2.0. | ||
|
||
package memcache | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
"sync" | ||
"time" | ||
|
||
"github.com/go-kit/kit/log" | ||
"github.com/go-kit/kit/log/level" | ||
"github.com/prometheus/client_golang/prometheus" | ||
"github.com/prometheus/client_golang/prometheus/promauto" | ||
"github.com/thanos-io/thanos/pkg/errutil" | ||
"github.com/thanos-io/thanos/pkg/extprom" | ||
) | ||
|
||
// Provider is a stateful cache for asynchronous memcached auto-discovery resolution. It provides a way to resolve | ||
// addresses and obtain them. | ||
type Provider struct { | ||
sync.RWMutex | ||
resolver Resolver | ||
clusterConfigs map[string]*clusterConfig | ||
logger log.Logger | ||
|
||
configVersion *extprom.TxGaugeVec | ||
resolvedAddresses *extprom.TxGaugeVec | ||
resolverFailuresCount prometheus.Counter | ||
resolverLookupsCount prometheus.Counter | ||
} | ||
|
||
func NewProvider(logger log.Logger, reg prometheus.Registerer, dialTimeout time.Duration) *Provider { | ||
p := &Provider{ | ||
resolver: &memcachedAutoDiscovery{dialTimeout: dialTimeout}, | ||
clusterConfigs: map[string]*clusterConfig{}, | ||
configVersion: extprom.NewTxGaugeVec(reg, prometheus.GaugeOpts{ | ||
Name: "auto_discovery_config_version", | ||
Help: "The current auto discovery config version", | ||
}, []string{"addr"}), | ||
resolvedAddresses: extprom.NewTxGaugeVec(reg, prometheus.GaugeOpts{ | ||
Name: "auto_discovery_resolved_addresses", | ||
Help: "The number of memcached nodes found via auto discovery", | ||
}, []string{"addr"}), | ||
resolverLookupsCount: promauto.With(reg).NewCounter(prometheus.CounterOpts{ | ||
Name: "auto_discovery_total", | ||
Help: "The number of memcache auto discovery attempts", | ||
}), | ||
resolverFailuresCount: promauto.With(reg).NewCounter(prometheus.CounterOpts{ | ||
Name: "auto_discovery_failures_total", | ||
Help: "The number of memcache auto discovery failures", | ||
}), | ||
logger: logger, | ||
} | ||
return p | ||
} | ||
|
||
// Resolve stores a list of nodes auto-discovered from the provided addresses. | ||
func (p *Provider) Resolve(ctx context.Context, addresses []string) error { | ||
clusterConfigs := map[string]*clusterConfig{} | ||
errs := errutil.MultiError{} | ||
|
||
for _, address := range addresses { | ||
clusterConfig, err := p.resolver.Resolve(ctx, address) | ||
p.resolverLookupsCount.Inc() | ||
|
||
if err != nil { | ||
level.Warn(p.logger).Log( | ||
"msg", "failed to perform auto-discovery for memcached", | ||
"address", address, | ||
) | ||
errs.Add(err) | ||
p.resolverFailuresCount.Inc() | ||
|
||
// Use cached values. | ||
p.RLock() | ||
clusterConfigs[address] = p.clusterConfigs[address] | ||
p.RUnlock() | ||
} else { | ||
clusterConfigs[address] = clusterConfig | ||
} | ||
} | ||
|
||
p.Lock() | ||
defer p.Unlock() | ||
|
||
p.resolvedAddresses.ResetTx() | ||
p.configVersion.ResetTx() | ||
for address, config := range clusterConfigs { | ||
p.resolvedAddresses.WithLabelValues(address).Set(float64(len(config.nodes))) | ||
p.configVersion.WithLabelValues(address).Set(float64(config.version)) | ||
} | ||
p.resolvedAddresses.Submit() | ||
p.configVersion.Submit() | ||
|
||
p.clusterConfigs = clusterConfigs | ||
|
||
return errs.Err() | ||
} | ||
|
||
// Addresses returns the latest addresses present in the Provider. | ||
func (p *Provider) Addresses() []string { | ||
p.RLock() | ||
defer p.RUnlock() | ||
|
||
var result []string | ||
for _, config := range p.clusterConfigs { | ||
for _, node := range config.nodes { | ||
result = append(result, fmt.Sprintf("%s:%d", node.dns, node.port)) | ||
} | ||
} | ||
return result | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,81 @@ | ||
// Copyright (c) The Thanos Authors. | ||
// Licensed under the Apache License 2.0. | ||
|
||
package memcache | ||
|
||
import ( | ||
"context" | ||
"sort" | ||
"testing" | ||
"time" | ||
|
||
"github.com/pkg/errors" | ||
|
||
"github.com/go-kit/kit/log" | ||
"github.com/thanos-io/thanos/pkg/testutil" | ||
) | ||
|
||
func TestProviderUpdatesAddresses(t *testing.T) { | ||
ctx := context.TODO() | ||
clusters := []string{"memcached-cluster-1", "memcached-cluster-2"} | ||
provider := NewProvider(log.NewNopLogger(), nil, 5*time.Second) | ||
resolver := mockResolver{ | ||
configs: map[string]*clusterConfig{ | ||
"memcached-cluster-1": {nodes: []node{{dns: "dns-1", ip: "ip-1", port: 11211}}}, | ||
"memcached-cluster-2": {nodes: []node{{dns: "dns-2", ip: "ip-2", port: 8080}}}, | ||
}, | ||
} | ||
provider.resolver = &resolver | ||
|
||
testutil.Ok(t, provider.Resolve(ctx, clusters)) | ||
addresses := provider.Addresses() | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080"}, addresses) | ||
|
||
resolver.configs = map[string]*clusterConfig{ | ||
"memcached-cluster-1": {nodes: []node{{dns: "dns-1", ip: "ip-1", port: 11211}, {dns: "dns-3", ip: "ip-3", port: 11211}}}, | ||
"memcached-cluster-2": {nodes: []node{{dns: "dns-2", ip: "ip-2", port: 8080}}}, | ||
} | ||
|
||
testutil.Ok(t, provider.Resolve(ctx, clusters)) | ||
addresses = provider.Addresses() | ||
sort.Strings(addresses) | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080", "dns-3:11211"}, addresses) | ||
} | ||
|
||
func TestProviderDoesNotUpdateAddressIfFailed(t *testing.T) { | ||
ctx := context.TODO() | ||
clusters := []string{"memcached-cluster-1", "memcached-cluster-2"} | ||
provider := NewProvider(log.NewNopLogger(), nil, 5*time.Second) | ||
resolver := mockResolver{ | ||
configs: map[string]*clusterConfig{ | ||
"memcached-cluster-1": {nodes: []node{{dns: "dns-1", ip: "ip-1", port: 11211}}}, | ||
"memcached-cluster-2": {nodes: []node{{dns: "dns-2", ip: "ip-2", port: 8080}}}, | ||
}, | ||
} | ||
provider.resolver = &resolver | ||
|
||
testutil.Ok(t, provider.Resolve(ctx, clusters)) | ||
addresses := provider.Addresses() | ||
sort.Strings(addresses) | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080"}, addresses) | ||
|
||
resolver.configs = nil | ||
resolver.err = errors.New("oops") | ||
|
||
testutil.NotOk(t, provider.Resolve(ctx, clusters)) | ||
addresses = provider.Addresses() | ||
sort.Strings(addresses) | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080"}, addresses) | ||
} | ||
|
||
type mockResolver struct { | ||
configs map[string]*clusterConfig | ||
err error | ||
} | ||
|
||
func (r *mockResolver) Resolve(_ context.Context, address string) (*clusterConfig, error) { | ||
if r.err != nil { | ||
return nil, r.err | ||
} | ||
return r.configs[address], nil | ||
} |
Oops, something went wrong.