tigerquoll commented on code in PR #1060:
URL: https://github.com/apache/yunikorn-k8shim/pull/1060#discussion_r3889398975
##########
pkg/admission/webhook_manager.go:
##########
@@ -86,17 +86,11 @@ type webhookManagerImpl struct {
// NewWebhookManager is used to create a new webhook manager
func NewWebhookManager(conf *conf.AdmissionControllerConf) (WebhookManager,
error) {
- kubeconfig, err := client.CreateRestConfig(conf.GetKubeConfig())
- if err != nil {
- log.Log(log.AdmissionWebhook).Error("Unable to create
kubernetes config", zap.Error(err))
- return nil, err
- }
- clientset, err := kubernetes.NewForConfig(kubeconfig)
- if err != nil {
- log.Log(log.AdmissionWebhook).Error("Unable to create
kubernetes clientset", zap.Error(err))
- return nil, err
- }
- return newWebhookManagerImpl(conf, clientset), nil
+ // the webhook and secret writes must be attributed to the admission
controller; the
+ // admission controller never populates the scheduler configuration, so
the client runs
Review Comment:
You're right on the outcome, and the comment was misleading either way, so
it's gone.
For the record the values didn't come from the configmap: nothing in the
admission controller calls `conf.UpdateConfigMaps`, so `NewKubeClient` read the
*defaults* of the scheduler configuration — 1000/1000 for the informer client
in `main.go` — while the webhook client, built straight from
`CreateRestConfig`, ran on client-go's 5/10. Both now run on the client-go
defaults, see the next thread.
##########
pkg/admission/webhook_manager.go:
##########
@@ -86,17 +86,11 @@ type webhookManagerImpl struct {
// NewWebhookManager is used to create a new webhook manager
func NewWebhookManager(conf *conf.AdmissionControllerConf) (WebhookManager,
error) {
- kubeconfig, err := client.CreateRestConfig(conf.GetKubeConfig())
- if err != nil {
- log.Log(log.AdmissionWebhook).Error("Unable to create
kubernetes config", zap.Error(err))
- return nil, err
- }
- clientset, err := kubernetes.NewForConfig(kubeconfig)
- if err != nil {
- log.Log(log.AdmissionWebhook).Error("Unable to create
kubernetes clientset", zap.Error(err))
- return nil, err
- }
- return newWebhookManagerImpl(conf, clientset), nil
+ // the webhook and secret writes must be attributed to the admission
controller; the
+ // admission controller never populates the scheduler configuration, so
the client runs
+ // on the defaults: no client side rate limiting
Review Comment:
Done. Both admission controller clients — the informers in `main.go` and the
webhook/secret writes in `NewWebhookManager` — now leave QPS and burst at 0,
i.e. client-go's 5/10, through `client.NewAdmissionControllerKubeClient`. The
bootstrap client is on the same defaults, which is what it had on master.
##########
pkg/client/apifactory.go:
##########
@@ -89,9 +89,18 @@ 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.
+func NewAPIFactory(scheduler api.SchedulerAPI, configs *conf.SchedulerConf,
testMode bool) (*APIFactory, error) {
kubeClient := NewKubeClient(configs.KubeConfig)
- namespaceInformerFactory :=
informers.NewSharedInformerFactoryWithOptions(kubeClient.GetClientSet(), 0,
informers.WithNamespace(configs.Namespace))
+
+ // all informers, cluster wide and namespaced, share one clientset:
both only run
+ // informers so they need the same unlimited client and the same
attribution.
Review Comment:
Agreed. With the events off the shared bucket, the relist-versus-event-storm
coupling that the informer client was meant to break is gone, and a separate
clientset only bought a distinct user agent in the audit log. The informer
clientset is removed; both informer factories are back on the scheduler client.
##########
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:
Done — the layout is now the two clients you describe: `yunikorn-scheduler`
for the scheduling actions and the informers, and `yunikorn-scheduler/events`
for the events, plus the pre-existing bootstrap and admission controller
clients, both on client-go defaults.
##########
pkg/client/eventsink.go:
##########
Review Comment:
Agreed that not all events are the same; the bucket was type-blind and that
was the wrong shape. What's changed:
- **Warning events are never shed by the bucket.** `FailedScheduling`,
`TaskRejected`, `NodeRejected`, `ApplicationFailed`, the gang scheduling
failures — everything that explains why a pod is not running — always goes out
and does not consume a token.
- **`PodUnschedulable` becomes a Warning.** It was typed Normal, which would
have put it in the sheddable class; kube-scheduler's equivalent,
`FailedScheduling`, is a Warning too.
- **Only Normal events are shed above `eventQPS`.** Counting what the shim
emits, the volume is three Normal events per *successfully* scheduled pod
(`Scheduling`, `Scheduled`, `PodBindSuccessful`) plus the core's
`Informational` records. At the ~1,900 binds/s measured on the KWOK rig that is
~5,700 events/s saying "this worked", and that is what a storm sheds.
- **`kubernetes.eventLevel`** (`normal` | `warning` | `none`, default
`normal`) is the verbosity knob you suggested. It is applied at the recorder,
so a suppressed event is never copied, cached or written, and it is
hot-reloadable like the log level.
The server-side mute still applies to every type. Once APF is rejecting
events, client-go drops them itself after its Retry-After retries
(`recordEvent` treats every `StatusError` as "Server rejected event (will not
retry!)"), so muting only saves the round trips and the parked goroutines; it
does not lose an event that would otherwise have landed.
##########
pkg/client/kubeclient.go:
##########
@@ -30,44 +30,127 @@ import (
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
"k8s.io/client-go/util/retry"
+ "k8s.io/utils/clock"
"github.com/apache/yunikorn-k8shim/pkg/conf"
"github.com/apache/yunikorn-k8shim/pkg/log"
)
type SchedulerKubeClient struct {
clientSet *kubernetes.Clientset
- configs *rest.Config
}
-func newBootstrapSchedulerKubeClient(kc string) SchedulerKubeClient {
- config := CreateRestConfigOrDie(kc)
- configuredClient, err := kubernetes.NewForConfig(config)
- if err != nil {
- log.Log(log.ShimClient).Fatal("failed to get Clientset",
zap.Error(err))
- }
- return SchedulerKubeClient{
- clientSet: configuredClient,
- configs: config,
+// Every client identifies the concern it serves in its user agent so that
traffic can be
+// attributed on the API server side. The concern comes first to allow prefix
matching, the
+// build version is appended by UserAgent().
+const (
+ UserAgentAdmissionController = "yunikorn-admission-controller"
+
+ userAgentBootstrap = "yunikorn-bootstrap"
+ userAgentWrites = "yunikorn-scheduler/writes"
+ userAgentInformers = "yunikorn-scheduler/informers"
+ userAgentEvents = "yunikorn-scheduler/events"
+)
+
+// UserAgent appends the build version to the concern, e.g.
"yunikorn-scheduler/writes
+// (1.7.0)". The version is only set in release builds, it is left out when
empty.
+func UserAgent(concern string) string {
+ return userAgentWithVersion(concern,
conf.GetBuildInfoMap()["buildVersion"])
Review Comment:
Removed; the user agents are the bare concern strings.
##########
pkg/client/kubeclient.go:
##########
@@ -30,44 +30,127 @@ import (
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
"k8s.io/client-go/util/retry"
+ "k8s.io/utils/clock"
"github.com/apache/yunikorn-k8shim/pkg/conf"
"github.com/apache/yunikorn-k8shim/pkg/log"
)
type SchedulerKubeClient struct {
clientSet *kubernetes.Clientset
- configs *rest.Config
}
-func newBootstrapSchedulerKubeClient(kc string) SchedulerKubeClient {
- config := CreateRestConfigOrDie(kc)
- configuredClient, err := kubernetes.NewForConfig(config)
- if err != nil {
- log.Log(log.ShimClient).Fatal("failed to get Clientset",
zap.Error(err))
- }
- return SchedulerKubeClient{
- clientSet: configuredClient,
- configs: config,
+// Every client identifies the concern it serves in its user agent so that
traffic can be
+// attributed on the API server side. The concern comes first to allow prefix
matching, the
+// build version is appended by UserAgent().
+const (
+ UserAgentAdmissionController = "yunikorn-admission-controller"
+
+ userAgentBootstrap = "yunikorn-bootstrap"
+ userAgentWrites = "yunikorn-scheduler/writes"
+ userAgentInformers = "yunikorn-scheduler/informers"
+ userAgentEvents = "yunikorn-scheduler/events"
+)
+
+// UserAgent appends the build version to the concern, e.g.
"yunikorn-scheduler/writes
+// (1.7.0)". The version is only set in release builds, it is left out when
empty.
+func UserAgent(concern string) string {
+ return userAgentWithVersion(concern,
conf.GetBuildInfoMap()["buildVersion"])
+}
+
+func userAgentWithVersion(concern, version string) string {
+ if version == "" {
+ return concern
}
+ return fmt.Sprintf("%s (%s)", concern, version)
}
-func newSchedulerKubeClient(kc string) SchedulerKubeClient {
- schedulerConf := conf.GetSchedulerConf()
+// rateLimitPolicy turns the configured qps and burst into the values to set
on a REST config
+// and a description of the resulting limit for logging.
+// A qps <= 0 disables client side rate limiting: this is expressed as a
negative QPS, which
+// client-go treats as "create no rate limiter", while a QPS of 0 silently
falls back to its
+// own defaults of 5 QPS / 10 burst (see RESTClientFor and
NewForConfigAndClient in
+// client-go). A QPS set without a burst is rejected by client-go, so the
burst defaults to
+// the configured qps.
Review Comment:
Done. `rateLimitPolicy` and its warnings are gone in favour of the
normalisation you suggested, with one difference: burst is the configured burst
when set, otherwise the qps — `max(burst, qps)` would silently raise an
explicit `burst < qps`. `0/0` is left to client-go's 5/10 as you asked, so
`qps: "0"` now behaves exactly as it did before this PR; the defaults stay at
`-1`. The compatibility table in the description is updated to match.
One thing I'd like to confirm with you: the scheduler client itself still
defaults to no client-side limiter (`-1`), on the measurements in the
description — the 1000/1000 bucket clipped the bind burst, and at `qps: 50` the
limiter was the throughput ceiling at 52 binds/s. If you'd rather keep a
default cap on the scheduler client, say so and I'll set one.
--
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]