From 44c9a3e94b6d4c9e78dadd20cd683755dec7523b Mon Sep 17 00:00:00 2001 From: xigang Date: Fri, 17 Jan 2020 18:48:12 +0800 Subject: [PATCH 1/2] fix etcd operator support workqueue(#1410) --- Gopkg.lock | 49 ++++++++++++++++ cmd/operator/main.go | 2 +- pkg/controller/controller.go | 7 +++ pkg/controller/informer.go | 101 +++++++++++++++++++++++---------- pkg/controller/sync.go | 106 +++++++++++++++++++++++++++++++++++ 5 files changed, 234 insertions(+), 31 deletions(-) create mode 100644 pkg/controller/sync.go diff --git a/Gopkg.lock b/Gopkg.lock index a346563b0..f5220e518 100644 --- a/Gopkg.lock +++ b/Gopkg.lock @@ -124,6 +124,7 @@ revision = "95a726a27e09030f9ccbd9982a1508f5a6d25ada" version = "v3.2.13" + [[projects]] digest = "1:4b8b5811da6970495e04d1f4e98bb89518cc3cfc3b3f456bdb876ed7b6c74049" name = "github.com/davecgh/go-spew" @@ -325,6 +326,17 @@ revision = "1624edc4454b8682399def8740d46db5e4362ba4" version = "v1.1.5" +[[projects]] + digest = "1:9aef72dc639ffc347ab8b81dc74d85ad12c64708e8bf8cdc43fac17b6018585d" + name = "github.com/kubernetes/client-go" + packages = [ + "tools/cache", + "util/workqueue", + ] + pruneopts = "NT" + revision = "78295b709ec6fa5be12e35892477a326dea2b5d3" + version = "kubernetes-1.12.6" + [[projects]] branch = "master" digest = "1:4925ec3736ef6c299cfcf61597782e3d66ec13114f7476019d04c742a7be55d0" @@ -987,10 +999,46 @@ "github.com/aws/aws-sdk-go/aws/session", "github.com/aws/aws-sdk-go/service/s3", "github.com/aws/aws-sdk-go/service/s3/s3manager", + "github.com/coreos/etcd-operator/pkg/apis/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/backup", + "github.com/coreos/etcd-operator/pkg/backup/backupapi", + "github.com/coreos/etcd-operator/pkg/backup/reader", + "github.com/coreos/etcd-operator/pkg/backup/util", + "github.com/coreos/etcd-operator/pkg/backup/writer", + "github.com/coreos/etcd-operator/pkg/chaos", + "github.com/coreos/etcd-operator/pkg/client", + "github.com/coreos/etcd-operator/pkg/cluster", + "github.com/coreos/etcd-operator/pkg/controller", + "github.com/coreos/etcd-operator/pkg/controller/backup-operator", + "github.com/coreos/etcd-operator/pkg/controller/restore-operator", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned/scheme", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned/typed/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned/typed/etcd/v1beta2/fake", + "github.com/coreos/etcd-operator/pkg/generated/informers/externalversions/etcd", + "github.com/coreos/etcd-operator/pkg/generated/informers/externalversions/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/generated/informers/externalversions/internalinterfaces", + "github.com/coreos/etcd-operator/pkg/generated/listers/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/util", + "github.com/coreos/etcd-operator/pkg/util/alibabacloudutil/ossfactory", + "github.com/coreos/etcd-operator/pkg/util/awsutil/s3factory", + "github.com/coreos/etcd-operator/pkg/util/azureutil/absfactory", + "github.com/coreos/etcd-operator/pkg/util/constants", + "github.com/coreos/etcd-operator/pkg/util/etcdutil", + "github.com/coreos/etcd-operator/pkg/util/gcputil/gcsfactory", + "github.com/coreos/etcd-operator/pkg/util/k8sutil", + "github.com/coreos/etcd-operator/pkg/util/probe", + "github.com/coreos/etcd-operator/pkg/util/retryutil", + "github.com/coreos/etcd-operator/test/e2e/e2eutil", + "github.com/coreos/etcd-operator/test/e2e/framework", + "github.com/coreos/etcd-operator/test/e2e/upgradetest/framework", + "github.com/coreos/etcd-operator/version", "github.com/coreos/etcd/clientv3", "github.com/coreos/etcd/etcdserver/api/v3rpc/rpctypes", "github.com/coreos/etcd/etcdserver/etcdserverpb", "github.com/coreos/etcd/pkg/transport", + "github.com/kubernetes/client-go/tools/cache", + "github.com/kubernetes/client-go/util/workqueue", "github.com/pborman/uuid", "github.com/pkg/errors", "github.com/prometheus/client_golang/prometheus", @@ -1004,6 +1052,7 @@ "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1beta1", "k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset", "k8s.io/apimachinery/pkg/api/errors", + "k8s.io/apimachinery/pkg/api/meta", "k8s.io/apimachinery/pkg/api/resource", "k8s.io/apimachinery/pkg/apis/meta/v1", "k8s.io/apimachinery/pkg/fields", diff --git a/cmd/operator/main.go b/cmd/operator/main.go index 54265610a..dc6e57d15 100644 --- a/cmd/operator/main.go +++ b/cmd/operator/main.go @@ -145,7 +145,7 @@ func run(ctx context.Context) { startChaos(context.Background(), cfg.KubeCli, cfg.Namespace, chaosLevel) c := controller.New(cfg) - err := c.Start() + err := c.Start(ctx) logrus.Fatalf("controller Start() failed: %v", err) } diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 8bdf1f49e..d98bed499 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -27,6 +27,8 @@ import ( apiextensionsclient "k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset" kwatch "k8s.io/apimachinery/pkg/watch" "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" ) var initRetryWaitTime = 30 * time.Second @@ -40,6 +42,11 @@ type Controller struct { logger *logrus.Entry Config + // k8s workqueue pattern + indexer cache.Indexer + informer cache.Controller + queue workqueue.RateLimitingInterface + clusters map[string]*cluster.Cluster } diff --git a/pkg/controller/informer.go b/pkg/controller/informer.go index 7d04b4cee..27f0435bc 100644 --- a/pkg/controller/informer.go +++ b/pkg/controller/informer.go @@ -25,8 +25,10 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" + "k8s.io/apimachinery/pkg/util/wait" kwatch "k8s.io/apimachinery/pkg/watch" "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" ) // TODO: get rid of this once we use workqueue @@ -36,7 +38,7 @@ func init() { pt = newPanicTimer(time.Minute, "unexpected long blocking (> 1 Minute) when handling cluster event") } -func (c *Controller) Start() error { +func (c *Controller) Start(ctx context.Context) error { // TODO: get rid of this init code. CRD and storage class will be managed outside of operator. for { err := c.initResource() @@ -49,11 +51,12 @@ func (c *Controller) Start() error { } probe.SetReady() - c.run() - panic("unreachable") + go c.run(ctx) + <-ctx.Done() + return ctx.Err() } -func (c *Controller) run() { +func (c *Controller) run(ctx context.Context) { var ns string if c.Config.ClusterWide { ns = metav1.NamespaceAll @@ -67,15 +70,29 @@ func (c *Controller) run() { ns, fields.Everything()) - _, informer := cache.NewIndexerInformer(source, &api.EtcdCluster{}, 0, cache.ResourceEventHandlerFuncs{ + c.queue = workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "etcd-operator") + c.indexer, c.informer = cache.NewIndexerInformer(source, &api.EtcdCluster{}, 0, cache.ResourceEventHandlerFuncs{ AddFunc: c.onAddEtcdClus, UpdateFunc: c.onUpdateEtcdClus, DeleteFunc: c.onDeleteEtcdClus, }, cache.Indexers{}) - ctx := context.TODO() - // TODO: use workqueue to avoid blocking - informer.Run(ctx.Done()) + defer c.queue.ShutDown() + + c.logger.Info("starting etcd controller") + go c.informer.Run(ctx.Done()) + + if !cache.WaitForCacheSync(ctx.Done(), c.informer.HasSynced) { + return + } + + const numWorkers = 1 + for i := 0; i < numWorkers; i++ { + go wait.Until(c.runWorker, time.Second, ctx.Done()) + } + + <-ctx.Done() + c.logger.Info("stopping etcd controller") } func (c *Controller) initResource() error { @@ -89,38 +106,62 @@ func (c *Controller) initResource() error { } func (c *Controller) onAddEtcdClus(obj interface{}) { - c.syncEtcdClus(obj.(*api.EtcdCluster)) + key, err := cache.MetaNamespaceKeyFunc(obj) + if err != nil { + panic(err) + } + c.queue.Add(key) } func (c *Controller) onUpdateEtcdClus(oldObj, newObj interface{}) { - c.syncEtcdClus(newObj.(*api.EtcdCluster)) + key, err := cache.MetaNamespaceKeyFunc(newObj) + if err != nil { + panic(err) + } + c.queue.Add(key) } func (c *Controller) onDeleteEtcdClus(obj interface{}) { - clus, ok := obj.(*api.EtcdCluster) - if !ok { - tombstone, ok := obj.(cache.DeletedFinalStateUnknown) - if !ok { - panic(fmt.Sprintf("unknown object from EtcdCluster delete event: %#v", obj)) - } - clus, ok = tombstone.Obj.(*api.EtcdCluster) - if !ok { - panic(fmt.Sprintf("Tombstone contained object that is not an EtcdCluster: %#v", obj)) - } - } - ev := &Event{ - Type: kwatch.Deleted, - Object: clus, - } - - pt.start() - _, err := c.handleClusterEvent(ev) + key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj) if err != nil { - c.logger.Warningf("fail to handle event: %v", err) + panic(err) } - pt.stop() + c.queue.Add(key) } +// func (c *Controller) onAddEtcdClus(obj interface{}) { +// c.syncEtcdClus(obj.(*api.EtcdCluster)) +// } + +// func (c *Controller) onUpdateEtcdClus(oldObj, newObj interface{}) { +// c.syncEtcdClus(newObj.(*api.EtcdCluster)) +// } + +// func (c *Controller) onDeleteEtcdClus(obj interface{}) { +// clus, ok := obj.(*api.EtcdCluster) +// if !ok { +// tombstone, ok := obj.(cache.DeletedFinalStateUnknown) +// if !ok { +// panic(fmt.Sprintf("unknown object from EtcdCluster delete event: %#v", obj)) +// } +// clus, ok = tombstone.Obj.(*api.EtcdCluster) +// if !ok { +// panic(fmt.Sprintf("Tombstone contained object that is not an EtcdCluster: %#v", obj)) +// } +// } +// ev := &Event{ +// Type: kwatch.Deleted, +// Object: clus, +// } + +// pt.start() +// _, err := c.handleClusterEvent(ev) +// if err != nil { +// c.logger.Warningf("fail to handle event: %v", err) +// } +// pt.stop() +// } + func (c *Controller) syncEtcdClus(clus *api.EtcdCluster) { ev := &Event{ Type: kwatch.Added, diff --git a/pkg/controller/sync.go b/pkg/controller/sync.go new file mode 100644 index 000000000..7e3744adb --- /dev/null +++ b/pkg/controller/sync.go @@ -0,0 +1,106 @@ +package controller + +import ( + api "github.com/coreos/etcd-operator/pkg/apis/etcd/v1beta2" + kwatch "k8s.io/apimachinery/pkg/watch" +) + +const ( + // Copy from deployment_controller.go: + // maxRetries is the number of times a etcd backup will be retried before it is dropped out of the queue. + // With the current rate-limiter in use (5ms*2^(maxRetries-1)) the following numbers represent the times + // an etcd backup is going to be requeued: + // + // 5ms, 10ms, 20ms, 40ms, 80ms, 160ms, 320ms, 640ms, 1.3s, 2.6s, 5.1s, 10.2s, 20.4s, 41s, 82s + maxRetries = 15 +) + +func (c *Controller) runWorker() { + for c.processNextItem() { + } +} + +func (c *Controller) processNextItem() bool { + // Wait until there is a new item in the working queue + key, quit := c.queue.Get() + if quit { + return false + } + + // Tell the queue that we are done with processing this key. This unblocks the key for other workers + // This allows safe parallel processing because two pods with the same key are never processed in + // parallel. + defer c.queue.Done(key) + err := c.processItem(key.(string)) + c.handleErr(err, key) + return true +} + +func (c *Controller) processItem(key string) error { + obj, exists, err := c.indexer.GetByKey(key) + if err != nil { + return err + } + + if !exists { + return nil + } + + clus := obj.(*api.EtcdCluster) + + ev := &Event{ + Type: kwatch.Added, + Object: clus, + } + + if clus.DeletionTimestamp != nil { + ev.Type = kwatch.Deleted + + pt.start() + _, err := c.handleClusterEvent(ev) + if err != nil { + c.logger.Warningf("fail to handle event: %v", err) + } + pt.stop() + } + + // re-watch or restart could give ADD event. + // If for an ADD event the cluster spec is invalid then it is not added to the local cache + // so modifying that cluster will result in another ADD event + if _, ok := c.clusters[getNamespacedName(clus)]; ok { + ev.Type = kwatch.Modified + } + + pt.start() + _, err = c.handleClusterEvent(ev) + if err != nil { + c.logger.Warningf("fail to handle event: %v", err) + return err + } + pt.stop() + return nil +} + +func (c *Controller) handleErr(err error, key interface{}) { + if err == nil { + // Forget about the #AddRateLimited history of the key on every successful synchronization. + // This ensures that future processing of updates for this key is not delayed because of + // an outdated error history. + c.queue.Forget(key) + return + } + + // This controller retries maxRetries times if something goes wrong. After that, it stops trying. + if c.queue.NumRequeues(key) < maxRetries { + c.logger.Errorf("error syncing etcd cluster (%v): %v", key, err) + + // Re-enqueue the key rate limited. Based on the rate limiter on the + // queue and the re-enqueue history, the key will be processed later again. + c.queue.AddRateLimited(key) + return + } + + c.queue.Forget(key) + // Report that, even after several retries, we could not successfully process this key + c.logger.Infof("Dropping etcd cluster(%v) out of the queue: %v", key, err) +} From 8d7f5f40a4e8cd648b987273ad1c5dea342a9c80 Mon Sep 17 00:00:00 2001 From: xigang Date: Sat, 18 Jan 2020 17:38:29 +0800 Subject: [PATCH 2/2] fix use workqueue to avoid blocking --- pkg/controller/informer.go | 54 -------------------------------------- pkg/controller/sync.go | 7 +++++ 2 files changed, 7 insertions(+), 54 deletions(-) diff --git a/pkg/controller/informer.go b/pkg/controller/informer.go index 27f0435bc..d812d86d1 100644 --- a/pkg/controller/informer.go +++ b/pkg/controller/informer.go @@ -26,7 +26,6 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/util/wait" - kwatch "k8s.io/apimachinery/pkg/watch" "k8s.io/client-go/tools/cache" "k8s.io/client-go/util/workqueue" ) @@ -129,59 +128,6 @@ func (c *Controller) onDeleteEtcdClus(obj interface{}) { c.queue.Add(key) } -// func (c *Controller) onAddEtcdClus(obj interface{}) { -// c.syncEtcdClus(obj.(*api.EtcdCluster)) -// } - -// func (c *Controller) onUpdateEtcdClus(oldObj, newObj interface{}) { -// c.syncEtcdClus(newObj.(*api.EtcdCluster)) -// } - -// func (c *Controller) onDeleteEtcdClus(obj interface{}) { -// clus, ok := obj.(*api.EtcdCluster) -// if !ok { -// tombstone, ok := obj.(cache.DeletedFinalStateUnknown) -// if !ok { -// panic(fmt.Sprintf("unknown object from EtcdCluster delete event: %#v", obj)) -// } -// clus, ok = tombstone.Obj.(*api.EtcdCluster) -// if !ok { -// panic(fmt.Sprintf("Tombstone contained object that is not an EtcdCluster: %#v", obj)) -// } -// } -// ev := &Event{ -// Type: kwatch.Deleted, -// Object: clus, -// } - -// pt.start() -// _, err := c.handleClusterEvent(ev) -// if err != nil { -// c.logger.Warningf("fail to handle event: %v", err) -// } -// pt.stop() -// } - -func (c *Controller) syncEtcdClus(clus *api.EtcdCluster) { - ev := &Event{ - Type: kwatch.Added, - Object: clus, - } - // re-watch or restart could give ADD event. - // If for an ADD event the cluster spec is invalid then it is not added to the local cache - // so modifying that cluster will result in another ADD event - if _, ok := c.clusters[getNamespacedName(clus)]; ok { - ev.Type = kwatch.Modified - } - - pt.start() - _, err := c.handleClusterEvent(ev) - if err != nil { - c.logger.Warningf("fail to handle event: %v", err) - } - pt.stop() -} - func (c *Controller) managed(clus *api.EtcdCluster) bool { if v, ok := clus.Annotations[k8sutil.AnnotationScope]; ok { if c.Config.ClusterWide { diff --git a/pkg/controller/sync.go b/pkg/controller/sync.go index 7e3744adb..f85b591d6 100644 --- a/pkg/controller/sync.go +++ b/pkg/controller/sync.go @@ -43,6 +43,13 @@ func (c *Controller) processItem(key string) error { } if !exists { + if _, ok := c.clusters[key]; !ok { + return nil + } + c.clusters[key].Delete() + delete(c.clusters, key) + clustersDeleted.Inc() + clustersTotal.Dec() return nil }