tigerquoll commented on code in PR #1060:
URL: https://github.com/apache/yunikorn-k8shim/pull/1060#discussion_r3766751286
##########
pkg/admission/webhook_manager.go:
##########
@@ -91,6 +91,11 @@ func NewWebhookManager(conf *conf.AdmissionControllerConf)
(WebhookManager, erro
log.Log(log.AdmissionWebhook).Error("Unable to create
kubernetes config", zap.Error(err))
return nil, err
}
+ // the webhook and secret writes must be attributed to the admission
controller
+ kubeconfig.UserAgent =
client.UserAgent(client.UserAgentAdmissionController)
+ // no client side rate limiting: leaving the QPS at 0 would mean the
client-go defaults
+ // of 5 QPS / 10 burst
+ kubeconfig.QPS = -1
Review Comment:
Yes, that works here too - done in the latest push: `NewWebhookManager` now
uses `NewKubeClientWithUserAgent` and takes `GetClientSet()` from it. The
policy comes out the same: the admission controller never populates the
scheduler configuration, so the constructor resolves to unlimited, which is
exactly what the hand-set `QPS = -1` did. Same user agent, and the webhook
client now gets the standard "creating Kubernetes client" log line the other
clients get.
##########
pkg/shim/scheduler.go:
##########
@@ -64,12 +64,12 @@ var (
)
func NewShimScheduler(scheduler api.SchedulerAPI, configs *conf.SchedulerConf,
bootstrapConfigMaps []*v1.ConfigMap) *KubernetesShim {
- kubeClient := client.NewKubeClient(configs.KubeConfig)
-
+ // all informers, cluster wide and namespaced, share one client
+ informerClientSet := client.NewInformerClientSet(configs.KubeConfig)
// we have disabled re-sync to keep ourselves up-to-date
- informerFactory :=
informers.NewSharedInformerFactory(kubeClient.GetClientSet(), 0)
+ informerFactory :=
informers.NewSharedInformerFactory(informerClientSet, 0)
Review Comment:
Nice idea - that made things even tidier. Done in the latest push: both
informer arguments are gone, the informer clientset and both factories are
created inside `NewAPIFactory`, and the signature keeps just `scheduler` and
`testMode` alongside `configs`.
##########
pkg/client/apifactory.go:
##########
@@ -89,9 +90,12 @@ type APIFactory struct {
lock *locking.RWMutex
}
-func NewAPIFactory(scheduler api.SchedulerAPI, informerFactory
informers.SharedInformerFactory, configs *conf.SchedulerConf, testMode bool)
(*APIFactory, error) {
+// NewAPIFactory creates the clients shared by the shim. The clientset backing
informerFactory
+// is passed in so that the namespaced factory created here shares it: both
only run informers
+// so they need the same unlimited client and the same attribution.
+func NewAPIFactory(scheduler api.SchedulerAPI, informerClientSet
kubernetes.Interface, informerFactory informers.SharedInformerFactory, configs
*conf.SchedulerConf, testMode bool) (*APIFactory, error) {
kubeClient := NewKubeClient(configs.KubeConfig)
Review Comment:
I think the writes client is the right one here, and it's what master
already does: this PR doesn't touch `LoadConfigMaps`, it just gives the client
it uses a name and a policy. There's no "reads" client to switch to - the
informers clientset exists only to back the informer factories and isn't
exposed in `Clients`, deliberately, so that list/watch traffic stays isolated
and attributable. The bootstrap client wouldn't be right either: it exists to
run before the configuration is loaded, so it builds from the default
kubeconfig path rather than the configured one, and its user agent is meant to
mark startup traffic only.
So the writes client carries the must-complete, direct request/response
traffic. The name is perhaps a little loose: apart from the configmap loads -
two GETs, once per process, at registration - everything on it is either a
mutation (pods, plus PVCs via the volume binder) or a read in service of one,
like the pod GETs inside the update retry loops. I've added a comment on
`Clients` in the latest push spelling out what each client carries.
##########
pkg/conf/schedulerconf.go:
##########
@@ -91,8 +93,10 @@ const (
DefaultOperatorPlugins = "general"
DefaultDisableGangScheduling = false
DefaultEnableConfigHotRefresh = true
- DefaultKubeQPS = 1000
- DefaultKubeBurst = 1000
+ DefaultKubeQPS = 0 // client side write
limiting is opt-in: <= 0 means no limiter
Review Comment:
Your understanding of client-go is right - a raw 0 in `rest.Config` lands on
the hidden 5/10 defaults. That's exactly why the configured value never reaches
the `rest.Config` as-is: `rateLimitPolicy` translates anything <= 0 into `QPS =
-1` before the client is built, which client-go documents as "disable
client-side ratelimiting". So the -1 you're suggesting is there - it's applied
at the point where client-go acts on it, and 0 in the configuration was just
the "not set" sentinel, with the documented convention being "<= 0 disables".
The same convention applies to `eventQPS`.
It's pinned by tests at each layer: `TestRateLimitPolicy` covers the 0 -> -1
translation, `TestNewRestConfig` asserts the resulting config carries QPS=-1
with no pre-built limiter, and `TestNewClientSetAcceptsRestConfig` runs every
shape we generate through `kubernetes.NewForConfig`. The kind run in the
description is the live confirmation - ~1,900 binds/s on defaults couldn't
happen under a 5/10 limiter.
That said, a careful reader taking the wrong meaning from the 0 is the best
argument that -1 communicates better - I've changed the defaults to -1 in the
latest push.
##########
pkg/client/eventsink.go:
##########
@@ -0,0 +1,228 @@
+/*
+ Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing, software
+ distributed under the License is distributed on an "AS IS" BASIS,
+ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ See the License for the specific language governing permissions and
+ limitations under the License.
+*/
+
+package client
+
+import (
+ "context"
+ "sync/atomic"
+ "time"
+
+ "go.uber.org/zap"
+ eventsv1 "k8s.io/api/events/v1"
+ "k8s.io/apimachinery/pkg/api/errors"
+ "k8s.io/client-go/tools/events"
+ "k8s.io/client-go/util/flowcontrol"
+ "k8s.io/utils/clock"
+
+ "github.com/apache/yunikorn-k8shim/pkg/conf"
+ "github.com/apache/yunikorn-k8shim/pkg/log"
+)
+
+const (
+ // the shortest time between two logs of the number of shed events
+ eventShedLogInterval = 30 * time.Second
+ // the longest the events are muted after the server asked us to back
off
+ maxEventMute = 60 * time.Second
+ // the mute used when the server throttles without saying for how long
+ defaultEventMute = time.Second
+)
+
+// rateLimitedEventSink sheds events which exceed the configured rate instead
of sending them.
+// The broadcaster writes every event from its own goroutine, so a rate
limiter on the events
+// client would only pace those writes: the goroutines pile up for as long as
a storm lasts.
+// Events are discardable, dropping them here bounds a storm instead of
delaying it.
+// The events are also muted for as long as the server asks us to back off:
priority and
+// fairness rejects with a 429 which carries the time to wait for.
+type rateLimitedEventSink struct {
+ inner events.EventSink
+ limiter flowcontrol.RateLimiter
+ rateLimit string
+ clock clock.PassiveClock
+ // the delay advertised by the server, recorded by the transport which
removes the header
+ hints *muteHintHolder
+
+ // number of events shed since they were last logged
+ shed atomic.Int64
+ // time the shed events were last logged, as unix nanoseconds
+ lastLogged atomic.Int64
+ // time the mute was last logged, as unix nanoseconds
+ lastMuteLogged atomic.Int64
+ // time until which no event is sent, as unix nanoseconds
+ muteUntil atomic.Int64
+}
+
+// NewRateLimitedEventSink wraps an event sink and sheds the events which
exceed the given
+// rate. A qps <= 0 disables shedding, every event is passed on to the wrapped
sink.
+func NewRateLimitedEventSink(inner events.EventSink, qps, burst int)
events.EventSink {
+ return newRateLimitedEventSink(inner, qps, burst, clock.RealClock{},
nil)
+}
+
+// newRateLimitedEventSink allows the clock and the delays recorded by the
transport to be
+// passed in for testing
+func newRateLimitedEventSink(inner events.EventSink, qps, burst int, clk
clock.PassiveClock, hints *muteHintHolder) *rateLimitedEventSink {
+ limitQPS, limitBurst, rateLimit := rateLimitPolicy(userAgentEvents,
qps, burst)
+
+ sink := &rateLimitedEventSink{
+ inner: inner,
+ rateLimit: rateLimit,
+ clock: clk,
+ hints: hints,
+ }
+ shedPolicy := rateLimit
+ if limitQPS > 0 {
+ sink.limiter = flowcontrol.NewTokenBucketRateLimiter(limitQPS,
limitBurst)
+ shedPolicy = "shed above " + rateLimit
+ }
+
+ log.Log(log.ShimClient).Info("creating event sink",
+ zap.String("concern", userAgentEvents),
+ zap.String("shedPolicy", shedPolicy))
+ return sink
+}
+
+// NewEventSink creates the sink for the event broadcaster: events are written
by the events
+// client and shed above the configured event rate.
+// The client and the sink share the delays advertised by the server: the
client fails fast on
+// a rejection instead of retrying it, the sink mutes the events until the
advertised deadline.
+func NewEventSink(kc string) events.EventSink {
+ schedulerConf := conf.GetSchedulerConf()
+
+ clk := clock.RealClock{}
+ hints := &muteHintHolder{}
+ sink := &events.EventSinkImpl{Interface: newEventsClientSet(kc, hints,
clk).EventsV1()}
+ return newRateLimitedEventSink(sink, schedulerConf.KubeEventQPS,
schedulerConf.KubeEventBurst, clk, hints)
+}
+
+func (s *rateLimitedEventSink) Create(ctx context.Context, event
*eventsv1.Event) (*eventsv1.Event, error) {
+ if s.shedEvent() {
+ return event, nil
+ }
+ result, err := s.inner.Create(ctx, event)
+ s.checkThrottled(err)
+ return result, err
+}
+
+func (s *rateLimitedEventSink) Update(ctx context.Context, event
*eventsv1.Event) (*eventsv1.Event, error) {
+ if s.shedEvent() {
+ return event, nil
+ }
+ result, err := s.inner.Update(ctx, event)
+ s.checkThrottled(err)
+ return result, err
+}
+
+func (s *rateLimitedEventSink) Patch(ctx context.Context, oldEvent
*eventsv1.Event, data []byte) (*eventsv1.Event, error) {
+ if s.shedEvent() {
+ return oldEvent, nil
+ }
+ result, err := s.inner.Patch(ctx, oldEvent, data)
+ s.checkThrottled(err)
+ return result, err
+}
+
+// checkThrottled mutes the events for as long as the server asked us to back
off. Priority
+// and fairness rejects with a 429 and advertises a delay which grows while it
keeps dropping
+// requests, following it is cheaper than having every event rejected.
+func (s *rateLimitedEventSink) checkThrottled(err error) {
+ if err == nil || !errors.IsTooManyRequests(err) {
+ return
+ }
+ delay, source := s.muteDelay(err)
+ if delay > maxEventMute {
+ delay = maxEventMute
+ }
+
+ muteUntil := s.clock.Now().Add(delay).UnixNano()
+ for {
+ current := s.muteUntil.Load()
+ // a mute is only ever extended, a later response must not
shorten it
+ if muteUntil <= current {
+ return
+ }
+ if s.muteUntil.CompareAndSwap(current, muteUntil) {
+ break
+ }
+ }
+
+ if s.claimLogInterval(&s.lastMuteLogged) {
+ log.Log(log.ShimClient).Warn("muting events, the server asked
to back off",
+ zap.Duration("delay", delay),
+ zap.String("delaySource", source),
+ zap.String("rateLimit", s.rateLimit))
+ }
+}
+
+// muteDelay returns how long the events must be muted and where that delay
came from. The
+// delay is normally recorded by the transport, which removes the header from
the response to
+// stop the REST client from retrying: the error the sink sees no longer
carries it. It is
+// still read from the error first, for the rejections which do not pass our
transport.
+func (s *rateLimitedEventSink) muteDelay(err error) (time.Duration, string) {
+ if seconds, ok := errors.SuggestsClientDelay(err); ok && seconds > 0 {
+ return time.Duration(seconds) * time.Second, "header"
+ }
+ if s.hints != nil {
+ if seconds, ok := s.hints.get(s.clock.Now()); ok && seconds > 0
{
+ return time.Duration(seconds) * time.Second, "transport"
+ }
+ }
+ return defaultEventMute, "default"
+}
+
+// shedEvent returns true if the event must be dropped instead of being passed
on. Dropping is
+// reported to the broadcaster as a successful write: it does not retry and
does not hold on
+// to the goroutine which is recording the event.
+// The mute is checked before the limiter: while the server is throttling us
there is no point
+// in spending a token on an event which it would reject.
+func (s *rateLimitedEventSink) shedEvent() bool {
+ muted := s.muted()
+ if !muted && (s.limiter == nil || s.limiter.TryAccept()) {
+ return false
+ }
+ s.shed.Add(1)
+ s.logShedEvents(muted)
+ return true
+}
+
+// muted returns true while the server asked us to stop sending events
+func (s *rateLimitedEventSink) muted() bool {
+ return s.clock.Now().UnixNano() < s.muteUntil.Load()
+}
+
+// logShedEvents logs the number of events shed since the last time they were
logged, at most
+// once every eventShedLogInterval. The first drop is always logged.
+func (s *rateLimitedEventSink) logShedEvents(muted bool) {
Review Comment:
The shed log can't use a rate-limited logger, unfortunately - it needs to
aggregate, not just throttle. Every dropped event increments a counter and the
call that claims the interval logs the accumulated total (`shed.Swap(0)`), so
the "9,805 events shed" line in the description is exact, not sampled. Since
shed events leave no other trace by design, that count is the only record of a
storm's size - it's the number you need when deciding whether `eventQPS` is set
right. The cadence of the two mechanisms is actually the same (first call logs,
then one per interval); the problem is the API shape: zap fields are evaluated
eagerly, so `Warn("events shed", zap.Int64("shedEvents", s.shed.Swap(0)))`
would execute the destructive `Swap(0)` before the logger decides to drop,
zeroing the count unlogged on every suppressed call. It also runs on the
injected clock, which is what lets the fake-clock tests cover the logging
cadence.
For context on why this client sheds rather than delays: events are the
lowest-priority traffic in the split - practically everything else is essential
for scheduling to actually occur - and the broadcaster spawns a goroutine per
event, so pacing them just converts an apiserver problem into unbounded
goroutine growth on our side. Queueing endlessly isn't an option and any static
queue size is never going to be the correct value, so we let the server tell us
"enough" via 429s and drop what it would refuse anyway. The fact and size of
every drop lands in the shim log.
Agreed a `RateLimitedLogger` would be a useful shim utility in general -
happy to file a follow-up for that.
##########
pkg/client/eventsink.go:
##########
@@ -0,0 +1,228 @@
+/*
+ Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing, software
+ distributed under the License is distributed on an "AS IS" BASIS,
+ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ See the License for the specific language governing permissions and
+ limitations under the License.
+*/
+
+package client
+
+import (
+ "context"
+ "sync/atomic"
+ "time"
+
+ "go.uber.org/zap"
+ eventsv1 "k8s.io/api/events/v1"
+ "k8s.io/apimachinery/pkg/api/errors"
+ "k8s.io/client-go/tools/events"
+ "k8s.io/client-go/util/flowcontrol"
+ "k8s.io/utils/clock"
+
+ "github.com/apache/yunikorn-k8shim/pkg/conf"
+ "github.com/apache/yunikorn-k8shim/pkg/log"
+)
+
+const (
+ // the shortest time between two logs of the number of shed events
+ eventShedLogInterval = 30 * time.Second
+ // the longest the events are muted after the server asked us to back
off
+ maxEventMute = 60 * time.Second
+ // the mute used when the server throttles without saying for how long
+ defaultEventMute = time.Second
+)
+
+// rateLimitedEventSink sheds events which exceed the configured rate instead
of sending them.
+// The broadcaster writes every event from its own goroutine, so a rate
limiter on the events
+// client would only pace those writes: the goroutines pile up for as long as
a storm lasts.
+// Events are discardable, dropping them here bounds a storm instead of
delaying it.
+// The events are also muted for as long as the server asks us to back off:
priority and
+// fairness rejects with a 429 which carries the time to wait for.
+type rateLimitedEventSink struct {
+ inner events.EventSink
+ limiter flowcontrol.RateLimiter
+ rateLimit string
+ clock clock.PassiveClock
+ // the delay advertised by the server, recorded by the transport which
removes the header
+ hints *muteHintHolder
+
+ // number of events shed since they were last logged
+ shed atomic.Int64
+ // time the shed events were last logged, as unix nanoseconds
+ lastLogged atomic.Int64
+ // time the mute was last logged, as unix nanoseconds
+ lastMuteLogged atomic.Int64
+ // time until which no event is sent, as unix nanoseconds
+ muteUntil atomic.Int64
+}
+
+// NewRateLimitedEventSink wraps an event sink and sheds the events which
exceed the given
+// rate. A qps <= 0 disables shedding, every event is passed on to the wrapped
sink.
+func NewRateLimitedEventSink(inner events.EventSink, qps, burst int)
events.EventSink {
+ return newRateLimitedEventSink(inner, qps, burst, clock.RealClock{},
nil)
+}
+
+// newRateLimitedEventSink allows the clock and the delays recorded by the
transport to be
+// passed in for testing
+func newRateLimitedEventSink(inner events.EventSink, qps, burst int, clk
clock.PassiveClock, hints *muteHintHolder) *rateLimitedEventSink {
+ limitQPS, limitBurst, rateLimit := rateLimitPolicy(userAgentEvents,
qps, burst)
+
+ sink := &rateLimitedEventSink{
+ inner: inner,
+ rateLimit: rateLimit,
+ clock: clk,
+ hints: hints,
+ }
+ shedPolicy := rateLimit
+ if limitQPS > 0 {
+ sink.limiter = flowcontrol.NewTokenBucketRateLimiter(limitQPS,
limitBurst)
+ shedPolicy = "shed above " + rateLimit
+ }
+
+ log.Log(log.ShimClient).Info("creating event sink",
+ zap.String("concern", userAgentEvents),
+ zap.String("shedPolicy", shedPolicy))
+ return sink
+}
+
+// NewEventSink creates the sink for the event broadcaster: events are written
by the events
+// client and shed above the configured event rate.
+// The client and the sink share the delays advertised by the server: the
client fails fast on
+// a rejection instead of retrying it, the sink mutes the events until the
advertised deadline.
+func NewEventSink(kc string) events.EventSink {
+ schedulerConf := conf.GetSchedulerConf()
+
+ clk := clock.RealClock{}
+ hints := &muteHintHolder{}
+ sink := &events.EventSinkImpl{Interface: newEventsClientSet(kc, hints,
clk).EventsV1()}
+ return newRateLimitedEventSink(sink, schedulerConf.KubeEventQPS,
schedulerConf.KubeEventBurst, clk, hints)
+}
+
+func (s *rateLimitedEventSink) Create(ctx context.Context, event
*eventsv1.Event) (*eventsv1.Event, error) {
+ if s.shedEvent() {
+ return event, nil
+ }
+ result, err := s.inner.Create(ctx, event)
+ s.checkThrottled(err)
+ return result, err
+}
+
+func (s *rateLimitedEventSink) Update(ctx context.Context, event
*eventsv1.Event) (*eventsv1.Event, error) {
+ if s.shedEvent() {
+ return event, nil
+ }
+ result, err := s.inner.Update(ctx, event)
+ s.checkThrottled(err)
+ return result, err
+}
+
+func (s *rateLimitedEventSink) Patch(ctx context.Context, oldEvent
*eventsv1.Event, data []byte) (*eventsv1.Event, error) {
+ if s.shedEvent() {
+ return oldEvent, nil
+ }
+ result, err := s.inner.Patch(ctx, oldEvent, data)
+ s.checkThrottled(err)
+ return result, err
+}
+
+// checkThrottled mutes the events for as long as the server asked us to back
off. Priority
+// and fairness rejects with a 429 and advertises a delay which grows while it
keeps dropping
+// requests, following it is cheaper than having every event rejected.
+func (s *rateLimitedEventSink) checkThrottled(err error) {
+ if err == nil || !errors.IsTooManyRequests(err) {
+ return
+ }
+ delay, source := s.muteDelay(err)
Review Comment:
`muteDelay` isn't re-detecting the 429 - the guard above does that - it
answers "for how long", and the hint can only answer that when it exists, is
fresh, and the holder is wired up. `muteDelay` covers the rest.
The apiserver's own filters always attach a numeric `Retry-After`, but the
sink can't assume every 429 on the wire came from them: a load balancer or
proxy in front of the apiserver can shed with a bare 429, in which case the
transport records nothing and the one-second default is what keeps that 429
from being a no-op. A hint that does exist expires after ten seconds, so a hint
from an earlier throttling episode can't set a fresh mute. And a sink built
without the transport at all - which is how `NewRateLimitedEventSink` and the
unit tests wire it - has no holder, so the delay can only be read from the
error itself; that's the first branch. The returned source also feeds the
`delaySource` field in the warning, which is how the APF test in the
description showed the transport path engaging.
Fetching the hint directly would still need the nil guard, the staleness
check and the fallback, so I'd rather keep them behind the one name.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]