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 <sanketsudake@gmail.com>
This commit is contained in:
Ankit Chawla
2022-04-05 21:51:05 +05:30
committed by GitHub
co-authored by Sanket Sudake
parent 4e01aa1322
commit 24dee03e18
7 changed files with 309 additions and 119 deletions
@@ -18,6 +18,10 @@ spec:
labels: labels:
svc: mqtrigger svc: mqtrigger
messagequeue: kafka messagequeue: kafka
annotations:
prometheus.io/scrape: "true"
prometheus.io/path: "/metrics"
prometheus.io/port: "8080"
spec: spec:
containers: containers:
- name: mqtrigger - name: mqtrigger
+3 -3
View File
@@ -65,8 +65,8 @@ func runTimer(ctx context.Context, logger *zap.Logger, routerUrl string) error {
return timer.Start(ctx, logger, routerUrl) return timer.Start(ctx, logger, routerUrl)
} }
func runMessageQueueMgr(logger *zap.Logger, routerUrl string) error { func runMessageQueueMgr(ctx context.Context, logger *zap.Logger, routerUrl string) error {
return mqtrigger.Start(logger, routerUrl) return mqtrigger.Start(ctx, logger, routerUrl)
} }
// KEDA based MessageQueue Trigger Manager // KEDA based MessageQueue Trigger Manager
@@ -276,7 +276,7 @@ Options:
} }
if arguments["--mqt"] == true { if arguments["--mqt"] == true {
err = runMessageQueueMgr(logger, routerUrl) err = runMessageQueueMgr(ctx, logger, routerUrl)
if err != nil { if err != nil {
logger.Error("message queue manager exited", zap.Error(err)) logger.Error("message queue manager exited", zap.Error(err))
return return
+4 -4
View File
@@ -17,6 +17,7 @@ limitations under the License.
package mqtrigger package mqtrigger
import ( import (
"context"
"fmt" "fmt"
"os" "os"
"path" "path"
@@ -35,7 +36,7 @@ import (
_ "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" _ "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() fissionClient, _, _, _, err := crd.MakeFissionClient()
if err != nil { if err != nil {
@@ -74,9 +75,8 @@ func Start(logger *zap.Logger, routerUrl string) error {
if err != nil { if err != nil {
logger.Fatal("failed to connect to remote message queue server", zap.Error(err)) logger.Fatal("failed to connect to remote message queue server", zap.Error(err))
} }
mqtMgr := mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, mq)
mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, mq).Run() mqtMgr.Run(ctx)
return nil return nil
} }
+23 -3
View File
@@ -33,6 +33,7 @@ import (
"go.uber.org/zap" "go.uber.org/zap"
fv1 "github.com/fission/fission/pkg/apis/core/v1" 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/factory"
"github.com/fission/fission/pkg/mqtrigger/messageQueue" "github.com/fission/fission/pkg/mqtrigger/messageQueue"
"github.com/fission/fission/pkg/mqtrigger/validator" "github.com/fission/fission/pkg/mqtrigger/validator"
@@ -74,6 +75,12 @@ type MqtConsumerGroupHandler struct {
fnUrl string fnUrl string
} }
type MqtConsumer struct {
ctx context.Context
cancel context.CancelFunc
consumer sarama.ConsumerGroup
}
func NewMqtConsumerGroupHandler(version sarama.KafkaVersion, func NewMqtConsumerGroupHandler(version sarama.KafkaVersion,
logger *zap.Logger, logger *zap.Logger,
trigger *fv1.MessageQueueTrigger, 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 { func (ch MqtConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for msg := range claim.Messages() { for msg := range claim.Messages() {
ch.kafkaMsgHandler(session, msg) ch.kafkaMsgHandler(session, msg)
mqtrigger.IncreaseMessageCount(ch.trigger.Name, ch.trigger.Namespace)
} }
return nil 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) ch := NewMqtConsumerGroupHandler(kafka.version, kafka.logger, trigger, producer, kafka.routerUrl)
// consume messages // consume messages
go func() { go func() {
topic := []string{trigger.Spec.Topic} topic := []string{trigger.Spec.Topic}
ctx := context.Background()
err = consumer.Consume(ctx, topic, ch) err = consumer.Consume(ctx, topic, ch)
if err != nil { if err != nil {
kafka.logger.Error("consumer error", zap.Error(err)) 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) { 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 { 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) { func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer sarama.SyncProducer, funcUrl string, err error, errorTopicHeaders []sarama.RecordHeader) {
+53
View File
@@ -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()
}
+124 -109
View File
@@ -19,21 +19,23 @@ package mqtrigger
import ( import (
"context" "context"
"errors" "errors"
"net/http"
"time" "time"
"github.com/prometheus/client_golang/prometheus/promhttp"
"go.uber.org/zap" "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" fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd" "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/mqtrigger/messageQueue"
"github.com/fission/fission/pkg/utils"
) )
const ( const (
ADD_TRIGGER requestType = iota ADD_TRIGGER requestType = iota
DELETE_TRIGGER DELETE_TRIGGER
GET_ALL_TRIGGERS GET_TRIGGER_SUBSCRIPTION
) )
type ( type (
@@ -59,8 +61,8 @@ type (
respChan chan response respChan chan response
} }
response struct { response struct {
err error err error
triggers *map[string]*triggerSubscription triggerSub *triggerSubscription
} }
) )
@@ -77,133 +79,146 @@ func MakeMessageQueueTriggerManager(logger *zap.Logger,
return &mqTriggerMgr return &mqTriggerMgr
} }
func (mqt *MessageQueueTriggerManager) Run() { func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) {
go mqt.service() 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() { func (mqt *MessageQueueTriggerManager) service() {
for { for {
req := <-mqt.reqChan 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 { switch req.requestType {
case ADD_TRIGGER: case ADD_TRIGGER:
var err error
k := crd.CacheKey(&req.triggerSub.trigger.ObjectMeta)
if _, ok := mqt.triggers[k]; ok { if _, ok := mqt.triggers[k]; ok {
err = errors.New("trigger already exists") resp.err = errors.New("trigger already exists")
} else { } else {
mqt.triggers[k] = req.triggerSub mqt.triggers[k] = req.triggerSub
mqt.logger.Debug("set trigger subscription", zap.String("key", k))
IncreaseSubscriptionCount()
} }
req.respChan <- response{err: err} req.respChan <- resp
case GET_ALL_TRIGGERS: case GET_TRIGGER_SUBSCRIPTION:
copyTriggers := make(map[string]*triggerSubscription) if _, ok := mqt.triggers[k]; !ok {
for key, val := range mqt.triggers { resp.err = errors.New("trigger does not exist")
copyTriggers[key] = val } else {
resp.triggerSub = mqt.triggers[k]
} }
req.respChan <- response{triggers: &copyTriggers} req.respChan <- resp
case DELETE_TRIGGER: 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 { func (mqt *MessageQueueTriggerManager) addTrigger(triggerSub *triggerSubscription) error {
respChan := make(chan response) resp := mqt.makeRequest(ADD_TRIGGER, triggerSub)
mqt.reqChan <- request{ return resp.err
requestType: ADD_TRIGGER,
triggerSub: triggerSub,
respChan: respChan,
}
r := <-respChan
return r.err
} }
func (mqt *MessageQueueTriggerManager) getAllTriggers() *map[string]*triggerSubscription { func (mqt *MessageQueueTriggerManager) getTriggerSubscription(trigger *fv1.MessageQueueTrigger) *triggerSubscription {
respChan := make(chan response) resp := mqt.makeRequest(GET_TRIGGER_SUBSCRIPTION, &triggerSubscription{trigger: *trigger})
mqt.reqChan <- request{ return resp.triggerSub
requestType: GET_ALL_TRIGGERS,
respChan: respChan,
}
r := <-respChan
return r.triggers
} }
func (mqt *MessageQueueTriggerManager) delTrigger(m *metav1.ObjectMeta) { func (mqt *MessageQueueTriggerManager) checkTriggerSubscription(trigger *fv1.MessageQueueTrigger) bool {
mqt.reqChan <- request{ return mqt.getTriggerSubscription(trigger) != nil
requestType: DELETE_TRIGGER, }
triggerSub: &triggerSubscription{
trigger: fv1.MessageQueueTrigger{ func (mqt *MessageQueueTriggerManager) delTriggerSubscription(trigger *fv1.MessageQueueTrigger) error {
ObjectMeta: *m, 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)
}
}
+98
View File
@@ -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")
}
}