From 13abb2e905a029cc89b9612800366f26a36f770f Mon Sep 17 00:00:00 2001 From: soharab-ic <156293296+soharab-ic@users.noreply.github.com> Date: Mon, 12 Aug 2024 18:00:41 +0530 Subject: [PATCH] Creating large number of MQTs takes time (#2984) * Implementing workqueue for MessageQueueTriggers * Fixing some issues with informers and deleteQueue * Fixing fission_mqt_created metrics * Rebase with main --------- Signed-off-by: Md Soharab Ansari --- cmd/fission-bundle/mqtrigger/mqtrigger.go | 20 +- pkg/mqtrigger/mqtmanager.go | 285 +++++++++++++++++----- pkg/mqtrigger/mqtmanager_test.go | 11 +- 3 files changed, 258 insertions(+), 58 deletions(-) diff --git a/cmd/fission-bundle/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go index aa3f5038..e45cbf32 100644 --- a/cmd/fission-bundle/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -22,16 +22,19 @@ import ( "os" "path" "strings" + "time" "github.com/pkg/errors" "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/mqtrigger" "github.com/fission/fission/pkg/mqtrigger/factory" "github.com/fission/fission/pkg/mqtrigger/messageQueue" _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" + "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/manager" ) @@ -73,11 +76,24 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * if err != nil { logger.Fatal("failed to connect to remote message queue server", zap.Error(err)) } - mqtMgr := mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, mq) - err = mqtMgr.Run(ctx, mgr) + + finformerFactory := make(map[string]genInformer.SharedInformerFactory, 0) + for _, ns := range utils.DefaultNSResolver().FissionResourceNS { + finformerFactory[ns] = genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil) + } + + mqtMgr, err := mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, finformerFactory, mq) if err != nil { return err } + + // Start informer factory + for _, factory := range finformerFactory { + factory.Start(ctx.Done()) + } + + mqtMgr.Run(ctx, ctx.Done(), mgr) + return nil } diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 37109019..619f9afe 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -19,15 +19,21 @@ package mqtrigger import ( "context" "errors" + "fmt" "time" "go.uber.org/zap" + apierrors "k8s.io/apimachinery/pkg/api/errors" + utilruntime "k8s.io/apimachinery/pkg/util/runtime" + "k8s.io/apimachinery/pkg/util/wait" k8sCache "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/generated/clientset/versioned" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" + flisterv1 "github.com/fission/fission/pkg/generated/listers/core/v1" "github.com/fission/fission/pkg/mqtrigger/messageQueue" - "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" ) @@ -48,6 +54,12 @@ type ( fissionClient versioned.Interface messageQueueType fv1.MessageQueueType messageQueue messageQueue.MessageQueue + + mqtLister map[string]flisterv1.MessageQueueTriggerLister + mqtListerSynced map[string]k8sCache.InformerSynced + + mqTriggerCreateUpdateQueue workqueue.RateLimitingInterface + mqTriggerDeleteQueue workqueue.RateLimitingInterface } triggerSubscription struct { @@ -67,36 +79,70 @@ type ( ) func MakeMessageQueueTriggerManager(logger *zap.Logger, - fissionClient versioned.Interface, mqType fv1.MessageQueueType, messageQueue messageQueue.MessageQueue) *MessageQueueTriggerManager { + fissionClient versioned.Interface, + mqType fv1.MessageQueueType, + finformerFactory map[string]genInformer.SharedInformerFactory, + messageQueue messageQueue.MessageQueue) (*MessageQueueTriggerManager, error) { mqTriggerMgr := MessageQueueTriggerManager{ - logger: logger.Named("message_queue_trigger_manager"), - reqChan: make(chan request), - triggers: make(map[string]*triggerSubscription), - fissionClient: fissionClient, - messageQueueType: mqType, - messageQueue: messageQueue, + logger: logger.Named("message_queue_trigger_manager"), + reqChan: make(chan request), + triggers: make(map[string]*triggerSubscription), + fissionClient: fissionClient, + mqtLister: make(map[string]flisterv1.MessageQueueTriggerLister, 0), + mqtListerSynced: make(map[string]k8sCache.InformerSynced, 0), + messageQueueType: mqType, + messageQueue: messageQueue, + mqTriggerCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "MqtAddUpdateQueue"), + mqTriggerDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "MqtDeleteQueue"), } - return &mqTriggerMgr + + for ns, informer := range finformerFactory { + _, err := informer.Core().V1().MessageQueueTriggers().Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: mqTriggerMgr.enqueueMqtAdd, + UpdateFunc: mqTriggerMgr.enqueueMqtUpdate, + DeleteFunc: mqTriggerMgr.enqueueMqtDelete, + }) + if err != nil { + return nil, err + } + mqTriggerMgr.mqtLister[ns] = informer.Core().V1().MessageQueueTriggers().Lister() + mqTriggerMgr.mqtListerSynced[ns] = informer.Core().V1().MessageQueueTriggers().Informer().HasSynced + } + + return &mqTriggerMgr, nil } -func (mqt *MessageQueueTriggerManager) Run(ctx context.Context, mgr manager.Interface) error { +func (mqt *MessageQueueTriggerManager) Run(ctx context.Context, stopCh <-chan struct{}, mgr manager.Interface) { + defer utilruntime.HandleCrash() + defer mqt.mqTriggerCreateUpdateQueue.ShutDown() + defer mqt.mqTriggerDeleteQueue.ShutDown() go mqt.service() - for _, informer := range utils.GetInformersForNamespaces(mqt.fissionClient, time.Minute*30, fv1.MessageQueueResource) { - _, err := informer.AddEventHandler(mqt.mqtInformerHandlers()) - if err != nil { - return err - } - mgr.Add(ctx, func(ctx context.Context) { - informer.Run(ctx.Done()) - }) - if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok { - mqt.logger.Fatal("failed to wait for caches to sync") - } + + mqt.logger.Info("Waiting for informer caches to sync") + + waitSynced := make([]k8sCache.InformerSynced, 0) + for _, synced := range mqt.mqtListerSynced { + waitSynced = append(waitSynced, synced) } + if ok := k8sCache.WaitForCacheSync(stopCh, waitSynced...); !ok { + mqt.logger.Fatal("failed to wait for caches to sync") + } + + for i := 0; i < 4; i++ { + mgr.Add(ctx, func(ctx context.Context) { + wait.Until(mqt.workerRun(ctx, "mqTriggerCreateUpdate", mqt.mqTriggerCreateUpdateQueueProcessFunc), time.Second, stopCh) + }) + } + mgr.Add(ctx, func(ctx context.Context) { + wait.Until(mqt.workerRun(ctx, "mqTriggerDeleteQueue", mqt.mqTriggerDeleteQueueProcessFunc), time.Second, stopCh) + }) + mgr.Add(ctx, func(ctx context.Context) { metrics.ServeMetrics(ctx, "mqtrigger", mqt.logger, mgr) }) - return nil + + <-stopCh + mqt.logger.Info("Shutting down workers for messageQueueTriggerManager") } func (mqt *MessageQueueTriggerManager) service() { @@ -161,22 +207,22 @@ func (mqt *MessageQueueTriggerManager) delTriggerSubscription(trigger *fv1.Messa return resp.err } -func (mqt *MessageQueueTriggerManager) RegisterTrigger(trigger *fv1.MessageQueueTrigger) { +func (mqt *MessageQueueTriggerManager) RegisterTrigger(trigger *fv1.MessageQueueTrigger) error { isPresent := mqt.checkTriggerSubscription(trigger) if isPresent { mqt.logger.Debug("message queue trigger already registered", zap.String("trigger_name", trigger.ObjectMeta.Name)) - return + return nil } // actually subscribe using the message queue client impl sub, err := mqt.messageQueue.Subscribe(trigger) if err != nil { mqt.logger.Warn("failed to subscribe to message queue trigger", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name)) - return + return err } if sub == nil { mqt.logger.Warn("subscription is nil", zap.String("trigger_name", trigger.ObjectMeta.Name)) - return + return nil } triggerSub := triggerSubscription{ trigger: *trigger, @@ -186,41 +232,170 @@ func (mqt *MessageQueueTriggerManager) RegisterTrigger(trigger *fv1.MessageQueue err = mqt.addTrigger(&triggerSub) if err != nil { mqt.logger.Fatal("adding message queue trigger failed", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name)) + return err } mqt.logger.Info("message queue trigger created", zap.String("trigger_name", trigger.ObjectMeta.Name)) + return nil } -func (mqt *MessageQueueTriggerManager) mqtInformerHandlers() k8sCache.ResourceEventHandlerFuncs { - return k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - trigger := obj.(*fv1.MessageQueueTrigger) - mqt.logger.Debug("Added mqt", zap.Any("trigger: ", trigger.ObjectMeta)) - mqt.RegisterTrigger(trigger) - }, - DeleteFunc: func(obj interface{}) { - trigger := obj.(*fv1.MessageQueueTrigger) - mqt.logger.Debug("Delete mqt", zap.Any("trigger: ", trigger.ObjectMeta)) - triggerSubscription := mqt.getTriggerSubscription(trigger) - if triggerSubscription == nil { - mqt.logger.Info("Unsubscribe failed", zap.String("trigger_name", trigger.ObjectMeta.Name)) - return - } +func (mqt *MessageQueueTriggerManager) enqueueMqtAdd(obj interface{}) { + key, err := k8sCache.MetaNamespaceKeyFunc(obj) + if err != nil { + mqt.logger.Error("error retrieving key from object in messageQueueTriggerManager", zap.Any("obj", obj)) + return + } + mqt.logger.Debug("enqueue mqt add", zap.String("key", key)) + mqt.mqTriggerCreateUpdateQueue.Add(key) +} - err := mqt.messageQueue.Unsubscribe(triggerSubscription.subscription) - if err != nil { - mqt.logger.Warn("failed to unsubscribe from message queue trigger", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name)) +func (mqt *MessageQueueTriggerManager) enqueueMqtUpdate(oldObj, newObj interface{}) { + key, err := k8sCache.MetaNamespaceKeyFunc(newObj) + if err != nil { + mqt.logger.Error("error retrieving key from object in messageQueueTriggerManager", zap.Any("obj", key)) + return + } + mqt.logger.Debug("enqueue mqt update", zap.String("key", key)) + mqt.mqTriggerCreateUpdateQueue.Add(key) +} + +func (mqt *MessageQueueTriggerManager) enqueueMqtDelete(obj interface{}) { + mqTrigger, ok := obj.(*fv1.MessageQueueTrigger) + if !ok { + mqt.logger.Error("unexpected type when deleting mqt to messageQueueTriggerManager", zap.Any("obj", obj)) + return + } + mqt.logger.Debug("enqueue mqt delete", zap.Any("mqTrigger", mqTrigger)) + mqt.mqTriggerDeleteQueue.Add(mqTrigger) +} + +func (mqt *MessageQueueTriggerManager) workerRun(ctx context.Context, name string, processFunc func(ctx context.Context) bool) func() { + return func() { + mqt.logger.Debug("Starting worker with func", zap.String("name", name)) + for { + if quit := processFunc(ctx); quit { + mqt.logger.Info("Shutting down worker", zap.String("name", name)) return } - err = mqt.delTriggerSubscription(trigger) - if err != nil { - mqt.logger.Warn("deleting message queue trigger failed", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name)) - } - mqt.logger.Info("message queue trigger deleted", zap.String("trigger_name", trigger.ObjectMeta.Name)) - }, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - trigger := newObj.(*fv1.MessageQueueTrigger) - mqt.logger.Debug("Updated mqt", zap.Any("trigger: ", trigger.ObjectMeta)) - mqt.RegisterTrigger(trigger) - }, + } } } + +func (mqt *MessageQueueTriggerManager) getMqtLister(namespace string) (flisterv1.MessageQueueTriggerLister, error) { + lister, ok := mqt.mqtLister[namespace] + if ok { + return lister, nil + } + for ns, lister := range mqt.mqtLister { + if ns == namespace { + return lister, nil + } + } + mqt.logger.Error("no messagequeuetrigger lister found for namespace", zap.String("namespace", namespace)) + return nil, fmt.Errorf("no messagequeuetrigger lister found for namespace %s", namespace) +} + +func (mqt *MessageQueueTriggerManager) mqTriggerCreateUpdateQueueProcessFunc(ctx context.Context) bool { + maxRetries := 3 + obj, quit := mqt.mqTriggerCreateUpdateQueue.Get() + if quit { + return false + } + key := obj.(string) + defer mqt.mqTriggerCreateUpdateQueue.Done(key) + + namespace, name, err := k8sCache.SplitMetaNamespaceKey(key) + if err != nil { + mqt.logger.Error("error splitting key", zap.Error(err)) + mqt.mqTriggerCreateUpdateQueue.Forget(key) + return false + } + mqTriggerLister, err := mqt.getMqtLister(namespace) + if err != nil { + mqt.logger.Error("error getting messagequeuetrigger lister", zap.Error(err)) + mqt.mqTriggerCreateUpdateQueue.Forget(key) + return false + } + mqTrigger, err := mqTriggerLister.MessageQueueTriggers(namespace).Get(name) + if apierrors.IsNotFound(err) { + mqt.logger.Info("mqt not found", zap.String("key", key)) + mqt.mqTriggerCreateUpdateQueue.Forget(key) + return false + } + + if err != nil { + if mqt.mqTriggerCreateUpdateQueue.NumRequeues(key) < maxRetries { + mqt.mqTriggerCreateUpdateQueue.AddRateLimited(key) + mqt.logger.Error("error getting mqt, retrying", zap.Error(err)) + } else { + mqt.mqTriggerCreateUpdateQueue.Forget(key) + mqt.logger.Error("error getting mqt, max retries reached", zap.Error(err)) + } + return false + } + + mqt.logger.Debug("Added mqt", zap.Any("trigger: ", mqTrigger.ObjectMeta)) + err = mqt.RegisterTrigger(mqTrigger) + if err != nil { + if mqt.mqTriggerCreateUpdateQueue.NumRequeues(key) < maxRetries { + mqt.mqTriggerCreateUpdateQueue.AddRateLimited(key) + mqt.logger.Error("error handling mqt from mqtInformer, retrying", zap.String("key", key), zap.Error(err)) + } else { + mqt.mqTriggerCreateUpdateQueue.Forget(key) + mqt.logger.Error("error handling mqt from mqtInformer, max retries reached", zap.String("key", key), zap.Error(err)) + } + return false + } + mqt.mqTriggerCreateUpdateQueue.Forget(key) + return false +} + +func (mqt *MessageQueueTriggerManager) mqTriggerDeleteQueueProcessFunc(ctx context.Context) bool { + maxRetries := 3 + obj, quit := mqt.mqTriggerDeleteQueue.Get() + if quit { + return false + } + defer mqt.mqTriggerDeleteQueue.Done(obj) + mqTrigger, ok := obj.(*fv1.MessageQueueTrigger) + if !ok { + mqt.logger.Error("unexpected type when deleting mqt to message queue trigger manager", zap.Any("obj", obj)) + mqt.mqTriggerDeleteQueue.Forget(obj) + return false + } + + mqt.logger.Debug("Delete mqt", zap.Any("trigger: ", mqTrigger.ObjectMeta)) + triggerSubscription := mqt.getTriggerSubscription(mqTrigger) + if triggerSubscription == nil { + mqt.logger.Info("Unsubscribe failed", zap.String("trigger_name", mqTrigger.ObjectMeta.Name)) + mqt.mqTriggerDeleteQueue.Forget(obj) + return false + } + + err := mqt.messageQueue.Unsubscribe(triggerSubscription.subscription) + if err != nil { + if mqt.mqTriggerDeleteQueue.NumRequeues(obj) < maxRetries { + mqt.mqTriggerDeleteQueue.AddRateLimited(obj) + mqt.logger.Error("failed to unsubscribe from message queue trigger, retrying", zap.Error(err), zap.String("trigger_name", mqTrigger.ObjectMeta.Name)) + } else { + mqt.mqTriggerDeleteQueue.Forget(obj) + mqt.logger.Error("failed to unsubscribe from message queue trigger, max retries reached", zap.Error(err)) + } + return false + } + + err = mqt.delTriggerSubscription(mqTrigger) + if err != nil { + if mqt.mqTriggerDeleteQueue.NumRequeues(obj) < maxRetries { + mqt.mqTriggerDeleteQueue.AddRateLimited(obj) + mqt.logger.Error("error deleting mqt, retrying", zap.Any("obj", obj), zap.Error(err)) + } else { + mqt.mqTriggerDeleteQueue.Forget(obj) + mqt.logger.Error("deleting message queue trigger failed, max retries reached", zap.Error(err), zap.String("trigger_name", mqTrigger.ObjectMeta.Name)) + } + return false + } + + mqt.mqTriggerDeleteQueue.Forget(obj) + mqt.logger.Info("message queue trigger deleted", zap.String("trigger_name", mqTrigger.ObjectMeta.Name)) + return false +} diff --git a/pkg/mqtrigger/mqtmanager_test.go b/pkg/mqtrigger/mqtmanager_test.go index 27025edc..d6d93969 100644 --- a/pkg/mqtrigger/mqtmanager_test.go +++ b/pkg/mqtrigger/mqtmanager_test.go @@ -19,10 +19,13 @@ package mqtrigger import ( "context" "testing" + "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" fv1 "github.com/fission/fission/pkg/apis/core/v1" + fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/mqtrigger/messageQueue" "github.com/fission/fission/pkg/utils/loggerfactory" ) @@ -54,7 +57,13 @@ func TestMqtManager(t *testing.T) { logger := loggerfactory.GetLogger() defer logger.Sync() msgQueue := fakeMessageQueue{} - mgr := MakeMessageQueueTriggerManager(logger, nil, fv1.MessageQueueTypeKafka, msgQueue) + fissionClient := fClient.NewSimpleClientset() + factory := make(map[string]genInformer.SharedInformerFactory, 0) + factory[metav1.NamespaceDefault] = genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, metav1.NamespaceDefault, nil) + mgr, err := MakeMessageQueueTriggerManager(logger, nil, fv1.MessageQueueTypeKafka, factory, msgQueue) + if err != nil { + t.Fatalf("Error creating messageQueueTriggerManagesr: %v", err) + } go mgr.service() trigger := fv1.MessageQueueTrigger{ ObjectMeta: metav1.ObjectMeta{