wilfred-s commented on code in PR #1060:
URL: https://github.com/apache/yunikorn-k8shim/pull/1060#discussion_r3854310023
##########
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:
The webhook should run with the default 5/10 or slightly higher which should
be more than enough. (3 informers and some adhoc connections on startup and
cert change)
All admission requests are REST calls made from the API server to the
controller and do not use a kubeclient using unlimited is not OK
##########
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:
That is not needed and also not correct. Informers do not do cycle through
connections or create new write/read requests on an ongoing basis. They only do
a single list on startup (1 connection) and then have a single watch on the
object (1 connection). The list and watch are successive and do not run at the
same time.
Watch is a push from the API on a long lived single API connection. There
are ~20 informers so we do not need unlimited. This client can be easily shared
with the normal scheduler actions without impact.
##########
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:
Incorrect it reads the configmap and loads the settings from `main()`. It
runs with the same QPS and burst as the scheduler. 1000/1000 which is too high
##########
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:
We should leave this to the K8s go-client and allow a QPS and burst set to 0
to use the defaults
Much simpler solution for the value check would be:
```
config.QPS = float32(max(-1, schedulerConf.KubeQPS))
config.Burst = max(max(0, schedulerConf.KubeBurst), schedulerConf.KubeQPS)
```
They come back from the config as int32 values
##########
pkg/client/eventsink.go:
##########
Review Comment:
I am not for dropping random events. Not all events are the same. Some you
really want to get pushed others are not that important. Limiting at the point
of creation for events that are of low impact is a better solution.
Even setting a verbosity level for events similar to the log level would be
acceptable.
##########
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:
Not sure if we want/should add build version to the agent string. It is not
required from a debugging perspective.
##########
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:
There should only be one split in the k8shim: events and scheduling actions.
The admission controller gets it own client already.
--
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]