From 24dee03e18716d0ca4131344610688b58ef222d7 Mon Sep 17 00:00:00 2001 From: Ankit Chawla Date: Tue, 5 Apr 2022 21:51:05 +0530 Subject: [PATCH] Added metrics for fission mqtrigger and optimizations in trigger subscriptions (#2399) - Use mqtrigger watch instead of polling mqtrigger every 5 seconds - Added metrics to monitor no of subscriptions, and no of messages per subscription - Add standard go metrics exported by prometheus - Enable prometheus discovery for mqtrigger pod - Optimized mqtrigger manager cache - Add unit tests for mqtrigger cache Co-authored-by: Sanket Sudake --- .../mqt-fission-kafka/deployment.yaml | 4 + cmd/fission-bundle/main.go | 6 +- cmd/fission-bundle/mqtrigger/mqtrigger.go | 8 +- pkg/mqtrigger/messageQueue/kafka/kafka.go | 26 +- pkg/mqtrigger/metrics.go | 53 ++++ pkg/mqtrigger/mqtmanager.go | 233 ++++++++++-------- pkg/mqtrigger/mqtmanager_test.go | 98 ++++++++ 7 files changed, 309 insertions(+), 119 deletions(-) create mode 100644 pkg/mqtrigger/metrics.go create mode 100644 pkg/mqtrigger/mqtmanager_test.go diff --git a/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml b/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml index a94ed0e7..47f959d6 100644 --- a/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml +++ b/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml @@ -18,6 +18,10 @@ spec: labels: svc: mqtrigger messagequeue: kafka + annotations: + prometheus.io/scrape: "true" + prometheus.io/path: "/metrics" + prometheus.io/port: "8080" spec: containers: - name: mqtrigger diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index 5f089594..50e7355c 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -65,8 +65,8 @@ func runTimer(ctx context.Context, logger *zap.Logger, routerUrl string) error { return timer.Start(ctx, logger, routerUrl) } -func runMessageQueueMgr(logger *zap.Logger, routerUrl string) error { - return mqtrigger.Start(logger, routerUrl) +func runMessageQueueMgr(ctx context.Context, logger *zap.Logger, routerUrl string) error { + return mqtrigger.Start(ctx, logger, routerUrl) } // KEDA based MessageQueue Trigger Manager @@ -276,7 +276,7 @@ Options: } if arguments["--mqt"] == true { - err = runMessageQueueMgr(logger, routerUrl) + err = runMessageQueueMgr(ctx, logger, routerUrl) if err != nil { logger.Error("message queue manager exited", zap.Error(err)) return diff --git a/cmd/fission-bundle/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go index 3628685c..cdea92a9 100644 --- a/cmd/fission-bundle/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -17,6 +17,7 @@ limitations under the License. package mqtrigger import ( + "context" "fmt" "os" "path" @@ -35,7 +36,7 @@ import ( _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" ) -func Start(logger *zap.Logger, routerUrl string) error { +func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { fissionClient, _, _, _, err := crd.MakeFissionClient() if err != nil { @@ -74,9 +75,8 @@ func Start(logger *zap.Logger, routerUrl string) error { if err != nil { logger.Fatal("failed to connect to remote message queue server", zap.Error(err)) } - - mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, mq).Run() - + mqtMgr := mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, mq) + mqtMgr.Run(ctx) return nil } diff --git a/pkg/mqtrigger/messageQueue/kafka/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go index 0ced9ac5..8cc8bf1b 100644 --- a/pkg/mqtrigger/messageQueue/kafka/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -33,6 +33,7 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "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/validator" @@ -74,6 +75,12 @@ type MqtConsumerGroupHandler struct { fnUrl string } +type MqtConsumer struct { + ctx context.Context + cancel context.CancelFunc + consumer sarama.ConsumerGroup +} + func NewMqtConsumerGroupHandler(version sarama.KafkaVersion, logger *zap.Logger, trigger *fv1.MessageQueueTrigger, @@ -114,6 +121,7 @@ func (ch MqtConsumerGroupHandler) Cleanup(sarama.ConsumerGroupSession) error { func (ch MqtConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { for msg := range claim.Messages() { ch.kafkaMsgHandler(session, msg) + mqtrigger.IncreaseMessageCount(ch.trigger.Name, ch.trigger.Namespace) } return nil } @@ -351,18 +359,28 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub } }() + ctx, cancel := context.WithCancel(context.Background()) ch := NewMqtConsumerGroupHandler(kafka.version, kafka.logger, trigger, producer, kafka.routerUrl) + // consume messages go func() { topic := []string{trigger.Spec.Topic} - ctx := context.Background() err = consumer.Consume(ctx, topic, ch) if err != nil { kafka.logger.Error("consumer error", zap.Error(err)) } + + if ctx.Err() != nil { + return + } }() - return consumer, nil + mqtConsumer := MqtConsumer{ + ctx: ctx, + cancel: cancel, + consumer: consumer, + } + return mqtConsumer, nil } func (kafka Kafka) getTLSConfig() (*tls.Config, error) { @@ -390,7 +408,9 @@ func (kafka Kafka) getTLSConfig() (*tls.Config, error) { } func (kafka Kafka) Unsubscribe(subscription messageQueue.Subscription) error { - return subscription.(sarama.ConsumerGroup).Close() + mqtConsumer := subscription.(MqtConsumer) + mqtConsumer.cancel() + return mqtConsumer.consumer.Close() } func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer sarama.SyncProducer, funcUrl string, err error, errorTopicHeaders []sarama.RecordHeader) { diff --git a/pkg/mqtrigger/metrics.go b/pkg/mqtrigger/metrics.go new file mode 100644 index 00000000..4bfdb282 --- /dev/null +++ b/pkg/mqtrigger/metrics.go @@ -0,0 +1,53 @@ +/* +Copyright 2022 The Fission Authors. + +Licensed 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 mqtrigger + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +var ( + metricsAddr = ":8080" + labels = []string{"trigger_name", "trigger_namespace"} + subscriptionCount = promauto.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "fission_mqt_subscriptions", + Help: "Total number of subscriptions to mq currently", + }, + []string{}, + ) + messageCount = promauto.NewCounterVec( + prometheus.CounterOpts{ + Name: "fission_mqt_messages_processed_total", + Help: "Total number of messages processed", + }, + labels, + ) +) + +func IncreaseSubscriptionCount() { + subscriptionCount.WithLabelValues().Inc() +} + +func DecreaseSubscriptionCount() { + subscriptionCount.WithLabelValues().Dec() +} + +func IncreaseMessageCount(trigname, trignamespace string) { + messageCount.WithLabelValues(trigname, trignamespace).Inc() +} diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 64ea08b8..2807f71c 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -19,21 +19,23 @@ package mqtrigger import ( "context" "errors" + "net/http" "time" + "github.com/prometheus/client_golang/prometheus/promhttp" "go.uber.org/zap" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8sCache "k8s.io/client-go/tools/cache" 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/messageQueue" - "github.com/fission/fission/pkg/utils" ) const ( ADD_TRIGGER requestType = iota DELETE_TRIGGER - GET_ALL_TRIGGERS + GET_TRIGGER_SUBSCRIPTION ) type ( @@ -59,8 +61,8 @@ type ( respChan chan response } response struct { - err error - triggers *map[string]*triggerSubscription + err error + triggerSub *triggerSubscription } ) @@ -77,133 +79,146 @@ func MakeMessageQueueTriggerManager(logger *zap.Logger, return &mqTriggerMgr } -func (mqt *MessageQueueTriggerManager) Run() { +func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) { go mqt.service() - go mqt.syncTriggers() + informerFactory := genInformer.NewSharedInformerFactory(mqt.fissionClient, time.Minute*30) + mqTriggerInformer := informerFactory.Core().V1().MessageQueueTriggers().Informer() + mqTriggerInformer.AddEventHandler(mqt.mqtInformerHandlers()) + go mqTriggerInformer.Run(ctx.Done()) + if ok := k8sCache.WaitForCacheSync(ctx.Done(), mqTriggerInformer.HasSynced); !ok { + mqt.logger.Fatal("failed to wait for caches to sync") + } + go mqt.serveMetrics() } func (mqt *MessageQueueTriggerManager) service() { for { req := <-mqt.reqChan + resp := response{triggerSub: nil, err: nil} + k, err := k8sCache.MetaNamespaceKeyFunc(&req.triggerSub.trigger) + if err != nil { + resp.err = err + req.respChan <- resp + continue + } + switch req.requestType { case ADD_TRIGGER: - var err error - k := crd.CacheKey(&req.triggerSub.trigger.ObjectMeta) if _, ok := mqt.triggers[k]; ok { - err = errors.New("trigger already exists") + resp.err = errors.New("trigger already exists") } else { mqt.triggers[k] = req.triggerSub + mqt.logger.Debug("set trigger subscription", zap.String("key", k)) + IncreaseSubscriptionCount() } - req.respChan <- response{err: err} - case GET_ALL_TRIGGERS: - copyTriggers := make(map[string]*triggerSubscription) - for key, val := range mqt.triggers { - copyTriggers[key] = val + req.respChan <- resp + case GET_TRIGGER_SUBSCRIPTION: + if _, ok := mqt.triggers[k]; !ok { + resp.err = errors.New("trigger does not exist") + } else { + resp.triggerSub = mqt.triggers[k] } - req.respChan <- response{triggers: ©Triggers} + req.respChan <- resp case DELETE_TRIGGER: - delete(mqt.triggers, crd.CacheKey(&req.triggerSub.trigger.ObjectMeta)) + delete(mqt.triggers, k) + mqt.logger.Debug("delete trigger", zap.String("key", k)) + DecreaseSubscriptionCount() + req.respChan <- resp } } } +func (mqt *MessageQueueTriggerManager) serveMetrics() { + http.Handle("/metrics", promhttp.Handler()) + err := http.ListenAndServe(metricsAddr, nil) + mqt.logger.Fatal("done listening on metrics endpoint", zap.Error(err)) +} + +func (mqt *MessageQueueTriggerManager) makeRequest(requestType requestType, triggerSub *triggerSubscription) response { + respChan := make(chan response) + mqt.reqChan <- request{requestType, triggerSub, respChan} + return <-respChan +} + func (mqt *MessageQueueTriggerManager) addTrigger(triggerSub *triggerSubscription) error { - respChan := make(chan response) - mqt.reqChan <- request{ - requestType: ADD_TRIGGER, - triggerSub: triggerSub, - respChan: respChan, - } - r := <-respChan - return r.err + resp := mqt.makeRequest(ADD_TRIGGER, triggerSub) + return resp.err } -func (mqt *MessageQueueTriggerManager) getAllTriggers() *map[string]*triggerSubscription { - respChan := make(chan response) - mqt.reqChan <- request{ - requestType: GET_ALL_TRIGGERS, - respChan: respChan, - } - r := <-respChan - return r.triggers +func (mqt *MessageQueueTriggerManager) getTriggerSubscription(trigger *fv1.MessageQueueTrigger) *triggerSubscription { + resp := mqt.makeRequest(GET_TRIGGER_SUBSCRIPTION, &triggerSubscription{trigger: *trigger}) + return resp.triggerSub } -func (mqt *MessageQueueTriggerManager) delTrigger(m *metav1.ObjectMeta) { - mqt.reqChan <- request{ - requestType: DELETE_TRIGGER, - triggerSub: &triggerSubscription{ - trigger: fv1.MessageQueueTrigger{ - ObjectMeta: *m, - }, +func (mqt *MessageQueueTriggerManager) checkTriggerSubscription(trigger *fv1.MessageQueueTrigger) bool { + return mqt.getTriggerSubscription(trigger) != nil +} + +func (mqt *MessageQueueTriggerManager) delTriggerSubscription(trigger *fv1.MessageQueueTrigger) error { + resp := mqt.makeRequest(DELETE_TRIGGER, &triggerSubscription{trigger: *trigger}) + return resp.err +} + +func (mqt *MessageQueueTriggerManager) RegisterTrigger(trigger *fv1.MessageQueueTrigger) { + isPresent := mqt.checkTriggerSubscription(trigger) + if isPresent { + mqt.logger.Info("message queue trigger already registered", zap.String("trigger_name", trigger.ObjectMeta.Name)) + return + } + + // 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 + } + if sub == nil { + mqt.logger.Warn("subscription is nil", zap.String("trigger_name", trigger.ObjectMeta.Name)) + return + } + triggerSub := triggerSubscription{ + trigger: *trigger, + subscription: sub, + } + // add to our list + 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)) + } + mqt.logger.Info("message queue trigger created", zap.String("trigger_name", trigger.ObjectMeta.Name)) +} + +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 + } + + 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)) + 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) syncTriggers() { - for { - // get new set of triggers - newTriggers, err := mqt.fissionClient.CoreV1().MessageQueueTriggers(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - if utils.IsNetworkError(err) { - mqt.logger.Error("encountered network error, will retry", zap.Error(err)) - time.Sleep(5 * time.Second) - continue - } - mqt.logger.Fatal("failed to read message queue trigger list", zap.Error(err)) - } - newTriggerMap := make(map[string]*fv1.MessageQueueTrigger) - for index := range newTriggers.Items { - newTrigger := &newTriggers.Items[index] - if newTrigger.Spec.MessageQueueType == mqt.messageQueueType { - newTriggerMap[crd.CacheKey(&newTrigger.ObjectMeta)] = newTrigger - } - } - - // get current set of triggers - currentTriggers := mqt.getAllTriggers() - - // register new triggers - for key, trigger := range newTriggerMap { - if _, ok := (*currentTriggers)[key]; ok { - continue - } - - // 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)) - continue - } - - triggerSub := triggerSubscription{ - trigger: *trigger, - subscription: sub, - } - - // add to our list - 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)) - } - - mqt.logger.Info("message queue trigger created", zap.String("trigger_name", trigger.ObjectMeta.Name)) - } - - // remove old triggers - for key, triggerSub := range *currentTriggers { - if _, ok := newTriggerMap[key]; ok { - continue - } - err := mqt.messageQueue.Unsubscribe(triggerSub.subscription) - if err != nil { - mqt.logger.Warn("failed to unsubscribe from message queue trigger", zap.Error(err), zap.String("trigger_name", triggerSub.trigger.ObjectMeta.Name)) - continue - } - mqt.delTrigger(&triggerSub.trigger.ObjectMeta) - mqt.logger.Info("message queue trigger deleted", zap.String("trigger_name", triggerSub.trigger.ObjectMeta.Name)) - } - - // TODO replace with a watch - time.Sleep(3 * time.Second) - } -} diff --git a/pkg/mqtrigger/mqtmanager_test.go b/pkg/mqtrigger/mqtmanager_test.go new file mode 100644 index 00000000..27025edc --- /dev/null +++ b/pkg/mqtrigger/mqtmanager_test.go @@ -0,0 +1,98 @@ +/* +Copyright 2022 The Fission Authors. + +Licensed 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 mqtrigger + +import ( + "context" + "testing" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/messageQueue" + "github.com/fission/fission/pkg/utils/loggerfactory" +) + +type mqtConsumer struct { + ctx context.Context + cancel context.CancelFunc +} + +type fakeMessageQueue struct { +} + +func (f fakeMessageQueue) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { + ctx, cancel := context.WithCancel(context.Background()) + mqtConsumer := mqtConsumer{ + ctx: ctx, + cancel: cancel, + } + return mqtConsumer, nil +} + +func (f fakeMessageQueue) Unsubscribe(triggerSub messageQueue.Subscription) error { + sub := triggerSub.(mqtConsumer) + sub.cancel() + return nil +} + +func TestMqtManager(t *testing.T) { + logger := loggerfactory.GetLogger() + defer logger.Sync() + msgQueue := fakeMessageQueue{} + mgr := MakeMessageQueueTriggerManager(logger, nil, fv1.MessageQueueTypeKafka, msgQueue) + go mgr.service() + trigger := fv1.MessageQueueTrigger{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test", + Namespace: "default", + }, + } + if mgr.checkTriggerSubscription(&trigger) { + t.Errorf("checkTrigger should return false") + } + sub, err := msgQueue.Subscribe(&trigger) + if err != nil { + t.Errorf("Subscribe should not return error") + } + triggerSub := triggerSubscription{ + trigger: trigger, + subscription: sub, + } + err = mgr.addTrigger(&triggerSub) + if err != nil { + t.Errorf("addTrigger should not return error") + } + if !mgr.checkTriggerSubscription(&trigger) { + t.Errorf("checkTrigger should return true") + } + getSub := mgr.getTriggerSubscription(&trigger) + if getSub == nil { + t.Fatal("getTriggerSubscription should return triggerSub") + } + if getSub.trigger.ObjectMeta.Name != trigger.ObjectMeta.Name { + t.Errorf("getTriggerSubscription should return triggerSub with trigger name %s", trigger.ObjectMeta.Name) + } + getSub.subscription.(mqtConsumer).cancel() + err = mgr.delTriggerSubscription(&trigger) + if err != nil { + t.Errorf("delTriggerSubscription should not return error") + } + if mgr.checkTriggerSubscription(&trigger) { + t.Errorf("checkTrigger should return false") + } +}