Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,19 @@ balancer.Register(consistent.NewBuilder(xxhash.Sum64))
grpc.Dial(addr, grpc.WithDefaultServiceConfig(consistent.DefaultServiceConfigJSON))
```

### Choosing an algorithm

By default, backends are placed on a consistent hashring with `ReplicationFactor` virtual nodes each.
Setting `Algorithm` to `consistent.AlgorithmRendezvous` uses rendezvous (highest random weight) hashing instead:

```go
cfg := &consistent.BalancerConfig{Spread: 1, Algorithm: consistent.AlgorithmRendezvous}
grpc.Dial(addr, grpc.WithDefaultServiceConfig(cfg.MustServiceConfigJSON()))
```

Rendezvous hashing spreads keys more evenly and makes membership changes far cheaper, at the cost of lookups that scale linearly with the number of backends.
See `go test ./rendezvous -run TestCompareDistribution -v` and `go test ./rendezvous -run '^$' -bench .` for a comparison.

## Acknowledgements

This project is a community effort fueled by contributions from both organizations and individuals.
Expand Down
153 changes: 153 additions & 0 deletions algorithm_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
package consistent

import (
"context"
"fmt"
"testing"

"github.com/cespare/xxhash/v2"
"github.com/stretchr/testify/require"
"google.golang.org/grpc/balancer"
"google.golang.org/grpc/connectivity"
"google.golang.org/grpc/resolver"

"github.com/authzed/consistent/hashring"
"github.com/authzed/consistent/rendezvous"
)

func TestParseConfigAlgorithm(t *testing.T) {
bld := NewBuilder(xxhash.Sum64)

cfg, err := bld.ParseConfig([]byte(`{"algorithm": "rendezvous"}`))
require.NoError(t, err)
require.Equal(t, &BalancerConfig{
ReplicationFactor: DefaultReplicationFactor,
Spread: DefaultSpread,
Algorithm: AlgorithmRendezvous,
}, cfg)

cfg, err = bld.ParseConfig([]byte(`{"algorithm": "hashring"}`))
require.NoError(t, err)
require.Equal(t, AlgorithmHashring, cfg.(*BalancerConfig).Algorithm)

_, err = bld.ParseConfig([]byte(`{"algorithm": "maglev"}`))
require.ErrorContains(t, err, `unknown algorithm "maglev"`)
}

func TestServiceConfigJSONAlgorithm(t *testing.T) {
got, err := (&BalancerConfig{Spread: 1, Algorithm: AlgorithmRendezvous}).ServiceConfigJSON()
require.NoError(t, err)
require.Equal(t, `{"loadBalancingConfig":[{"consistent-hashring":{"spread":1,"algorithm":"rendezvous"}}]}`, got)
}

// rendezvousBalancer builds a balancer that uses rendezvous hashing and
// moves every SubConn for addrs to READY.
func rendezvousBalancer(t *testing.T, addrs ...resolver.Address) *ringBalancer {
t.Helper()
b, _ := readyBalancer(t, addrs...)
require.NoError(t, b.UpdateClientConnState(balancer.ClientConnState{
ResolverState: resolver.State{Addresses: addrs},
BalancerConfig: &BalancerConfig{ReplicationFactor: 100, Spread: 1, Algorithm: AlgorithmRendezvous},
}))
require.IsType(t, &rendezvous.Set{}, b.hashring)
return b
}

// Switching algorithms rebuilds the member set. The new set must contain
// every member except those in TRANSIENT_FAILURE, and switching back must
// restore a hashring.
func TestAlgorithmChangeRebuildsMemberSet(t *testing.T) {
addrs := []resolver.Address{
{ServerName: "t", Addr: "1"},
{ServerName: "t", Addr: "2"},
{ServerName: "t", Addr: "3"},
}
b, _ := readyBalancer(t, addrs...)
require.IsType(t, &hashring.Ring{}, b.hashring)

sci2, _ := b.subConns.Get(addrs[1])
b.UpdateSubConnState(sci2.(balancer.SubConn), balancer.SubConnState{
ConnectivityState: connectivity.TransientFailure,
ConnectionError: fmt.Errorf("refused"),
})

require.NoError(t, b.UpdateClientConnState(balancer.ClientConnState{
ResolverState: resolver.State{Addresses: addrs},
BalancerConfig: &BalancerConfig{ReplicationFactor: 100, Spread: 1, Algorithm: AlgorithmRendezvous},
}))
require.IsType(t, &rendezvous.Set{}, b.hashring)
require.ElementsMatch(t, []string{"t1", "t3"}, ringKeys(b))

require.NoError(t, b.UpdateClientConnState(balancer.ClientConnState{
ResolverState: resolver.State{Addresses: addrs},
BalancerConfig: &BalancerConfig{ReplicationFactor: 100, Spread: 1},
}))
require.IsType(t, &hashring.Ring{}, b.hashring)
require.ElementsMatch(t, []string{"t1", "t3"}, ringKeys(b))
}

// Rendezvous hashing has no virtual nodes, so a ReplicationFactor change must
// not rebuild the member set and move keys.
func TestRendezvousIgnoresReplicationFactor(t *testing.T) {
addrs := []resolver.Address{{ServerName: "t", Addr: "1"}, {ServerName: "t", Addr: "2"}}
b := rendezvousBalancer(t, addrs...)
before := b.hashring

require.NoError(t, b.UpdateClientConnState(balancer.ClientConnState{
ResolverState: resolver.State{Addresses: addrs},
BalancerConfig: &BalancerConfig{ReplicationFactor: 7, Spread: 1, Algorithm: AlgorithmRendezvous},
}))
require.Same(t, before, b.hashring)
}

// With rendezvous hashing, the picker, the RingView, and a standalone
// rendezvous.Set with the same members all agree, and a failed backend's
// keys move to a remaining one.
func TestRendezvousPickerRoutes(t *testing.T) {
bld := NewBuilder(xxhash.Sum64)
addrs := []resolver.Address{{ServerName: "t", Addr: "1"}, {ServerName: "t", Addr: "2"}}
b := readyBalancerForTarget(t, bld, "test:///backends", addrs...)
require.NoError(t, b.UpdateClientConnState(balancer.ClientConnState{
ResolverState: resolver.State{Addresses: addrs},
BalancerConfig: &BalancerConfig{ReplicationFactor: 100, Spread: 1, Algorithm: AlgorithmRendezvous},
}))
view := bld.RingFor("test:///backends")

ref := rendezvous.New(xxhash.Sum64)
for _, a := range addrs {
require.NoError(t, ref.Add(subConnMember{key: a.ServerName + a.Addr}))
}

pick := func(key []byte) balancer.SubConn {
res, err := b.picker.Pick(balancer.PickInfo{Ctx: context.WithValue(context.Background(), CtxKey, key)})
require.NoError(t, err)
return res.SubConn
}

var keyFor2 []byte
for i := 0; i < 100; i++ {
key := []byte(fmt.Sprintf("key-%d", i))
want, err := ref.FindN(key, 1)
require.NoError(t, err)

viewed, err := view.FindN(key, 1)
require.NoError(t, err)
require.Equal(t, want[0].Key(), viewed[0].Key())

picked := pick(key)
require.Equal(t, want[0].Key(), b.scKeys[picked])

if want[0].Key() == "t2" {
keyFor2 = key
}
}
require.NotNil(t, keyFor2, "no key hashes to t2")

sci2, _ := b.subConns.Get(addrs[1])
b.UpdateSubConnState(sci2.(balancer.SubConn), balancer.SubConnState{
ConnectivityState: connectivity.TransientFailure,
ConnectionError: fmt.Errorf("refused"),
})
sci1, _ := b.subConns.Get(addrs[0])
require.Equal(t, sci1.(balancer.SubConn), pick(keyFor2))
}
71 changes: 65 additions & 6 deletions balancer.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import (
"google.golang.org/grpc/serviceconfig"

"github.com/authzed/consistent/hashring"
"github.com/authzed/consistent/rendezvous"
)

type ctxKey string
Expand All @@ -49,6 +50,14 @@ const (
// DefaultSpread is the value that will be used when parsing a service
// config provides an invalid value.
DefaultSpread = 1

// AlgorithmHashring places backends on a consistent hashring with
// ReplicationFactor virtual nodes each. This is the default.
AlgorithmHashring = "hashring"

// AlgorithmRendezvous selects backends with rendezvous (highest random
// weight) hashing. It ignores ReplicationFactor.
AlgorithmRendezvous = "rendezvous"
)

// DefaultServiceConfigJSON is a helper to easily leverage the defaults.
Expand All @@ -73,6 +82,46 @@ type BalancerConfig struct {
serviceconfig.LoadBalancingConfig `json:"-"`
ReplicationFactor uint16 `json:"replicationFactor,omitempty"`
Spread uint8 `json:"spread,omitempty"`
// Algorithm is AlgorithmHashring or AlgorithmRendezvous. Empty means
// AlgorithmHashring.
Algorithm string `json:"algorithm,omitempty"`
}

func (c *BalancerConfig) algorithm() string {
if c.Algorithm == "" {
return AlgorithmHashring
}
return c.Algorithm
}

// needsNewMemberSet reports whether moving from the old config to the new one
// requires building a new member set.
func needsNewMemberSet(old, updated *BalancerConfig) bool {
if old == nil || old.algorithm() != updated.algorithm() {
return true
}
return updated.algorithm() == AlgorithmHashring && old.ReplicationFactor != updated.ReplicationFactor
}

// memberSet is the structure a balancer places its backends on. Both
// hashring.Ring and rendezvous.Set implement it.
type memberSet interface {
Add(hashring.Member) error
Remove(hashring.Member) error
FindN(key []byte, num uint8) ([]hashring.Member, error)
Members() []hashring.Member
}

var (
_ memberSet = (*hashring.Ring)(nil)
_ memberSet = (*rendezvous.Set)(nil)
)

func newMemberSet(hashfn hashring.HashFunc, cfg *BalancerConfig) memberSet {
if cfg.algorithm() == AlgorithmRendezvous {
return rendezvous.New(hashfn)
}
return hashring.MustNew(hashfn, cfg.ReplicationFactor)
}

// ServiceConfigJSON encodes the current config into the gRPC Service Config
Expand Down Expand Up @@ -211,6 +260,12 @@ func (b *builder) ParseConfig(js json.RawMessage) (serviceconfig.LoadBalancingCo
lbCfg.Spread = DefaultSpread
}

switch lbCfg.Algorithm {
case "", AlgorithmHashring, AlgorithmRendezvous:
default:
return nil, fmt.Errorf("consistent-hashring: unknown algorithm %q", lbCfg.Algorithm)
}

return &lbCfg, nil
}

Expand All @@ -228,11 +283,14 @@ type ringBalancer struct {
ringMembers map[balancer.SubConn]struct{}

config *BalancerConfig
hashring *hashring.Ring
hashring memberSet
hasher hashring.HashFunc
// slot is where the balancer publishes its live ring for RingView
// readers.
slot *ringSlot
// published is the entry this balancer stored in slot, so that Close
// withdraws only its own ring.
published *publishedRing

resolverErr error // the last error reported by the resolver; cleared on successful resolution
connErr error // the last connection error; cleared upon leaving TransientFailure
Expand Down Expand Up @@ -274,9 +332,10 @@ func (b *ringBalancer) UpdateClientConnState(s balancer.ClientConnState) error {
// update the service config if it has changed
if s.BalancerConfig != nil {
svcConfig := s.BalancerConfig.(*BalancerConfig)
if b.config == nil || svcConfig.ReplicationFactor != b.config.ReplicationFactor {
b.hashring = hashring.MustNew(b.hasher, svcConfig.ReplicationFactor)
b.slot.ring.Store(b.hashring)
if needsNewMemberSet(b.config, svcConfig) {
b.hashring = newMemberSet(b.hasher, svcConfig)
b.published = &publishedRing{b.hashring}
b.slot.ring.Store(b.published)
// The new ring starts empty: put every SubConn that has not
// failed back on it.
b.ringMembers = make(map[balancer.SubConn]struct{})
Expand Down Expand Up @@ -488,7 +547,7 @@ func (b *ringBalancer) Close() {
// balancer that no longer exists. The CompareAndSwap clears only this
// balancer's own ring, so it cannot remove the ring of a second balancer
// that shares the target's slot.
b.slot.ring.CompareAndSwap(b.hashring, nil)
b.slot.ring.CompareAndSwap(b.published, nil)
}

func (b *ringBalancer) ExitIdle() {
Expand All @@ -498,7 +557,7 @@ func (b *ringBalancer) ExitIdle() {
}

type picker struct {
hashring *hashring.Ring
hashring memberSet
spread uint8
}

Expand Down
Loading
Loading