improve concurrent watcher registration perf

This commit is contained in:
Forrest Marshall
2021-12-09 13:01:35 -08:00
committed by Forrest
parent d52241d969
commit dfd3732c6b
3 changed files with 180 additions and 12 deletions
+2 -2
View File
@@ -350,7 +350,7 @@ type Cache struct {
webSessionCache types.WebSessionInterface
webTokenCache types.WebTokenInterface
windowsDesktopsCache services.WindowsDesktops
eventsFanout *services.Fanout
eventsFanout *services.FanoutSet
// closed indicates that the cache has been closed
closed *atomic.Bool
@@ -624,7 +624,7 @@ func New(config Config) (*Cache, error) {
webSessionCache: local.NewIdentityService(wrapper).WebSessions(),
webTokenCache: local.NewIdentityService(wrapper).WebTokens(),
windowsDesktopsCache: local.NewWindowsDesktopService(wrapper),
eventsFanout: services.NewFanout(),
eventsFanout: services.NewFanoutSet(),
Entry: log.WithFields(log.Fields{
trace.Component: config.Component,
}),
+107 -10
View File
@@ -23,6 +23,8 @@ import (
"github.com/gravitational/teleport/api/types"
"github.com/gravitational/trace"
"go.uber.org/atomic"
)
const defaultQueueSize = 64
@@ -201,7 +203,7 @@ func (f *Fanout) Emit(events ...types.Event) {
func (f *Fanout) Reset() {
f.mu.Lock()
defer f.mu.Unlock()
f.closeWatchers()
f.closeWatchersAsync()
f.init = false
}
@@ -210,19 +212,31 @@ func (f *Fanout) Reset() {
func (f *Fanout) Close() {
f.mu.Lock()
defer f.mu.Unlock()
f.closeWatchers()
f.closeWatchersAsync()
f.closed = true
}
func (f *Fanout) closeWatchers() {
for _, entries := range f.watchers {
for _, entry := range entries {
entry.watcher.cancel()
}
}
// watcher map was potentially quite large, so
// relenguish that memory.
// closeWatchersAsync moves ownership of the watcher mapping to a background goroutine
// for asynchronous cancellation and sets up a new empty mapping.
func (f *Fanout) closeWatchersAsync() {
watchersToClose := f.watchers
f.watchers = make(map[string][]fanoutEntry)
// goroutines run with a "happens after" releationship to the
// expressions that create them. since we move ownership of the
// old watcher mapping prior to spawning this goroutine, we are
// "safe" to modify it without worrying about locking. because
// we don't continue to hold the lock in the foreground goroutine,
// this fanout instance may permit new events/registrations/inits/resets
// while the old watchers are still being closed. this is fine, since
// the aformentioned move guarantees that these old watchers aren't
// going to observe any of the new state transitions.
go func() {
for _, entries := range watchersToClose {
for _, entry := range entries {
entry.watcher.cancel()
}
}
}()
}
func (f *Fanout) addWatcher(w *fanoutWatcher) {
@@ -354,3 +368,86 @@ func (w *fanoutWatcher) Error() error {
return nil
}
}
// fanoutSetSize is the number of members in a fanout set. selected based on some experimentation with
// the FanoutSetRegistration benchmark. This value keeps 100K concurrent registrations well under 1s.
const fanoutSetSize = 128
// FanoutSet is a collection of separate Fanout instances. It exposes an identical API, and "load balances"
// watcher registration across the enclosed instances. In very large clusters it is possible for tens of
// thousands of nodes to simultaneously request watchers. This can cause serious contention issues. FanoutSet is
// a simple but effective solution to that problem.
type FanoutSet struct {
// rw mutex is used to ensure that Close and Reset operations are exclusive,
// since these operations close watchers. Enforcing this property isn't strictly
// necessary, but it prevents a scenario where watchers might observe a reset/close,
// attempt re-registration, and observe the *same* reset/close again. This isn't
// necessarily a problem, but it might confuse attempts to debug other event-system
// issues, so we choose to avoid it.
rw sync.RWMutex
counter *atomic.Uint64
members []*Fanout
}
// NewFanoutSet creates a new FanoutSet instance in an uninitialized
// state. Until initialized, watchers will be queued but no
// events will be sent.
func NewFanoutSet() *FanoutSet {
members := make([]*Fanout, 0, fanoutSetSize)
for i := 0; i < fanoutSetSize; i++ {
members = append(members, NewFanout())
}
return &FanoutSet{
counter: atomic.NewUint64(0),
members: members,
}
}
// NewWatcher attaches a new watcher to a fanout instance.
func (s *FanoutSet) NewWatcher(ctx context.Context, watch types.Watch) (types.Watcher, error) {
s.rw.RLock() // see field-level docks for locking model
defer s.rw.RUnlock()
fi := int(s.counter.Inc() % uint64(len(s.members)))
return s.members[fi].NewWatcher(ctx, watch)
}
// SetInit sets the Fanout instances into an initialized state, sending OpInit
// events to any watchers which were added prior to initialization.
func (s *FanoutSet) SetInit() {
s.rw.RLock() // see field-level docks for locking model
defer s.rw.RUnlock()
for _, f := range s.members {
f.SetInit()
}
}
// Emit broadcasts events to all matching watchers that have been attached
// to this fanout set.
func (s *FanoutSet) Emit(events ...types.Event) {
s.rw.RLock() // see field-level docks for locking model
defer s.rw.RUnlock()
for _, f := range s.members {
f.Emit(events...)
}
}
// Reset closes all attached watchers and places the fanout instances
// into an uninitialized state. Reset may be called on an uninitialized
// fanout set to remove "queued" watchers.
func (s *FanoutSet) Reset() {
s.rw.Lock() // see field-level docks for locking model
defer s.rw.Unlock()
for _, f := range s.members {
f.Reset()
}
}
// Close permanently closes the fanout. Existing watchers will be
// closed and no new watchers will be added.
func (s *FanoutSet) Close() {
s.rw.Lock() // see field-level docks for locking model
defer s.rw.Unlock()
for _, f := range s.members {
f.Close()
}
}
+71
View File
@@ -18,6 +18,7 @@ package services
import (
"context"
"sync"
"testing"
"time"
@@ -70,3 +71,73 @@ func TestFanoutInit(t *testing.T) {
default:
}
}
/*
goos: linux
goarch: amd64
pkg: github.com/gravitational/teleport/lib/services
cpu: Intel(R) Core(TM) i9-10885H CPU @ 2.40GHz
BenchmarkFanoutRegistration-16 1 118856478045 ns/op
*/
// NOTE: this benchmark exists primarily to "contrast" with the set registration
// benchmark below, and demonstrate why the set-based strategy is necessary.
func BenchmarkFanoutRegistration(b *testing.B) {
const iterations = 100_000
ctx := context.Background()
for n := 0; n < b.N; n++ {
f := NewFanout()
f.SetInit()
var wg sync.WaitGroup
for i := 0; i < iterations; i++ {
wg.Add(1)
go func() {
defer wg.Done()
w, err := f.NewWatcher(ctx, types.Watch{
Name: "test",
Kinds: []types.WatchKind{{Name: "spam"}, {Name: "eggs"}},
})
require.NoError(b, err)
w.Close()
}()
}
wg.Wait()
}
}
/*
goos: linux
goarch: amd64
pkg: github.com/gravitational/teleport/lib/services
cpu: Intel(R) Core(TM) i9-10885H CPU @ 2.40GHz
BenchmarkFanoutSetRegistration-16 3 394211563 ns/op
*/
func BenchmarkFanoutSetRegistration(b *testing.B) {
const iterations = 100_000
ctx := context.Background()
for n := 0; n < b.N; n++ {
f := NewFanoutSet()
f.SetInit()
var wg sync.WaitGroup
for i := 0; i < iterations; i++ {
wg.Add(1)
go func() {
defer wg.Done()
w, err := f.NewWatcher(ctx, types.Watch{
Name: "test",
Kinds: []types.WatchKind{{Name: "spam"}, {Name: "eggs"}},
})
require.NoError(b, err)
w.Close()
}()
}
wg.Wait()
}
}