From 90d781ca2d0d45b51f6fea2d1efcac9261090a6d Mon Sep 17 00:00:00 2001 From: soharab-ic <156293296+soharab-ic@users.noreply.github.com> Date: Fri, 9 Aug 2024 11:58:44 +0530 Subject: [PATCH] Fixed mqtrigger scaling issue (#2986) * Fixed mqtrigger scaling issue Signed-off-by: Md Soharab Ansari * Resolve review comments Signed-off-by: Md Soharab Ansari * Set goreleaser version to v1 Remove armv7 references from goreleaser file Signed-off-by: Md Soharab Ansari --------- Signed-off-by: Md Soharab Ansari --- .github/workflows/release.yaml | 2 +- .goreleaser.yml | 104 ------------------- pkg/mqtrigger/messageQueue/kafka/consumer.go | 4 + pkg/mqtrigger/messageQueue/kafka/kafka.go | 49 ++++----- pkg/mqtrigger/metrics.go | 32 ++++++ 5 files changed, 58 insertions(+), 133 deletions(-) diff --git a/.github/workflows/release.yaml b/.github/workflows/release.yaml index 7219b0e0..db8ac181 100644 --- a/.github/workflows/release.yaml +++ b/.github/workflows/release.yaml @@ -82,7 +82,7 @@ jobs: - name: Run GoReleaser uses: goreleaser/goreleaser-action@7ec5c2b0c6cdda6e8bbb49444bc797dd33d74dd8 # v5.0.0 with: - version: latest + version: "~> v1" args: release env: COSIGN_PWD: ${{ secrets.COSIGN_PWD }} diff --git a/.goreleaser.yml b/.goreleaser.yml index 17ff697c..bdaf14b9 100644 --- a/.goreleaser.yml +++ b/.goreleaser.yml @@ -241,191 +241,87 @@ dockers: - "--label=org.opencontainers.image.created={{.Date}}" - "--label=org.opencontainers.image.revision={{.FullCommit}}" - "--label=org.opencontainers.image.version={{.Tag}}" - - &docker-armv7 - use: buildx - goos: linux - goarch: arm - goarm: 7 - ids: - - builder - image_templates: - - "fission/builder:latest-armv7" - - "fission/builder:{{ .Tag }}-armv7" - - "ghcr.io/fission/builder:latest-armv7" - - "ghcr.io/fission/builder:{{ .Tag }}-armv7" - dockerfile: cmd/builder/Dockerfile - build_flag_templates: - - "--label=org.opencontainers.image.description=The builder assists in building the fission function source code for deployment." - - "--label=org.opencontainers.image.source={{.GitURL}}" - - "--platform=linux/arm/v7" - - "--label=org.opencontainers.image.created={{.Date}}" - - "--label=org.opencontainers.image.revision={{.FullCommit}}" - - "--label=org.opencontainers.image.version={{.Tag}}" - - <<: *docker-armv7 - ids: - - fetcher - image_templates: - - "fission/fetcher:latest-armv7" - - "fission/fetcher:{{ .Tag }}-armv7" - - "ghcr.io/fission/fetcher:latest-armv7" - - "ghcr.io/fission/fetcher:{{ .Tag }}-armv7" - dockerfile: cmd/fetcher/Dockerfile - build_flag_templates: - - "--label=org.opencontainers.image.description=Fetcher is a lightweight component used by environment and builder pods. Fetcher helps in fetch and upload of source/deployment packages and specializing environments." - - "--label=org.opencontainers.image.source={{.GitURL}}" - - "--platform=linux/arm/v7" - - "--label=org.opencontainers.image.created={{.Date}}" - - "--label=org.opencontainers.image.revision={{.FullCommit}}" - - "--label=org.opencontainers.image.version={{.Tag}}" - - <<: *docker-armv7 - ids: - - fission-bundle - image_templates: - - "fission/fission-bundle:latest-armv7" - - "fission/fission-bundle:{{ .Tag }}-armv7" - - "ghcr.io/fission/fission-bundle:latest-armv7" - - "ghcr.io/fission/fission-bundle:{{ .Tag }}-armv7" - dockerfile: cmd/fission-bundle/Dockerfile - build_flag_templates: - - "--label=org.opencontainers.image.description=fission-bundle is a component which is a single binary for all components. Most server side components running on server side are fission-bundle binary wrapped in container and used with different arguments." - - "--label=org.opencontainers.image.source={{.GitURL}}" - - "--platform=linux/arm/v7" - - "--label=org.opencontainers.image.created={{.Date}}" - - "--label=org.opencontainers.image.revision={{.FullCommit}}" - - "--label=org.opencontainers.image.version={{.Tag}}" - - <<: *docker-armv7 - ids: - - pre-upgrade-checks - image_templates: - - "fission/pre-upgrade-checks:latest-armv7" - - "fission/pre-upgrade-checks:{{ .Tag }}-armv7" - - "ghcr.io/fission/pre-upgrade-checks:latest-armv7" - - "ghcr.io/fission/pre-upgrade-checks:{{ .Tag }}-armv7" - dockerfile: cmd/preupgradechecks/Dockerfile - build_flag_templates: - - "--label=org.opencontainers.image.description=Preupgradechecks ensures that Fission is ready for the targeted version upgrade by performing checks beforehand." - - "--label=org.opencontainers.image.source={{.GitURL}}" - - "--platform=linux/arm/v7" - - "--label=org.opencontainers.image.created={{.Date}}" - - "--label=org.opencontainers.image.revision={{.FullCommit}}" - - "--label=org.opencontainers.image.version={{.Tag}}" - - <<: *docker-armv7 - ids: - - reporter - image_templates: - - "fission/reporter:latest-armv7" - - "fission/reporter:{{ .Tag }}-armv7" - - "ghcr.io/fission/reporter:latest-armv7" - - "ghcr.io/fission/reporter:{{ .Tag }}-armv7" - dockerfile: cmd/reporter/Dockerfile - build_flag_templates: - - "--label=org.opencontainers.image.description=The reporter gathers information that assists in improving fission." - - "--label=org.opencontainers.image.source={{.GitURL}}" - - "--platform=linux/arm/v7" - - "--label=org.opencontainers.image.created={{.Date}}" - - "--label=org.opencontainers.image.revision={{.FullCommit}}" - - "--label=org.opencontainers.image.version={{.Tag}}" docker_manifests: - name_template: ghcr.io/fission/builder:{{ .Tag }} image_templates: - ghcr.io/fission/builder:{{ .Tag }}-amd64 - ghcr.io/fission/builder:{{ .Tag }}-arm64 - - ghcr.io/fission/builder:{{ .Tag }}-armv7 - name_template: fission/builder:{{ .Tag }} image_templates: - fission/builder:{{ .Tag }}-amd64 - fission/builder:{{ .Tag }}-arm64 - - fission/builder:{{ .Tag }}-armv7 - name_template: ghcr.io/fission/fetcher:{{ .Tag }} image_templates: - ghcr.io/fission/fetcher:{{ .Tag }}-amd64 - ghcr.io/fission/fetcher:{{ .Tag }}-arm64 - - ghcr.io/fission/fetcher:{{ .Tag }}-armv7 - name_template: fission/fetcher:{{ .Tag }} image_templates: - fission/fetcher:{{ .Tag }}-amd64 - fission/fetcher:{{ .Tag }}-arm64 - - fission/fetcher:{{ .Tag }}-armv7 - name_template: ghcr.io/fission/fission-bundle:{{ .Tag }} image_templates: - ghcr.io/fission/fission-bundle:{{ .Tag }}-amd64 - ghcr.io/fission/fission-bundle:{{ .Tag }}-arm64 - - ghcr.io/fission/fission-bundle:{{ .Tag }}-armv7 - name_template: fission/fission-bundle:{{ .Tag }} image_templates: - fission/fission-bundle:{{ .Tag }}-amd64 - fission/fission-bundle:{{ .Tag }}-arm64 - - fission/fission-bundle:{{ .Tag }}-armv7 - name_template: ghcr.io/fission/pre-upgrade-checks:{{ .Tag }} image_templates: - ghcr.io/fission/pre-upgrade-checks:{{ .Tag }}-amd64 - ghcr.io/fission/pre-upgrade-checks:{{ .Tag }}-arm64 - - ghcr.io/fission/pre-upgrade-checks:{{ .Tag }}-armv7 - name_template: fission/pre-upgrade-checks:{{ .Tag }} image_templates: - fission/pre-upgrade-checks:{{ .Tag }}-amd64 - fission/pre-upgrade-checks:{{ .Tag }}-arm64 - - fission/pre-upgrade-checks:{{ .Tag }}-armv7 - name_template: ghcr.io/fission/reporter:{{ .Tag }} image_templates: - ghcr.io/fission/reporter:{{ .Tag }}-amd64 - ghcr.io/fission/reporter:{{ .Tag }}-arm64 - - ghcr.io/fission/reporter:{{ .Tag }}-armv7 - name_template: fission/reporter:{{ .Tag }} image_templates: - fission/reporter:{{ .Tag }}-amd64 - fission/reporter:{{ .Tag }}-arm64 - - fission/reporter:{{ .Tag }}-armv7 - name_template: ghcr.io/fission/builder:latest image_templates: - ghcr.io/fission/builder:latest-amd64 - ghcr.io/fission/builder:latest-arm64 - - ghcr.io/fission/builder:latest-armv7 - name_template: fission/builder:latest image_templates: - fission/builder:latest-amd64 - fission/builder:latest-arm64 - - fission/builder:latest-armv7 - name_template: ghcr.io/fission/fetcher:latest image_templates: - ghcr.io/fission/fetcher:latest-amd64 - ghcr.io/fission/fetcher:latest-arm64 - - ghcr.io/fission/fetcher:latest-armv7 - name_template: fission/fetcher:latest image_templates: - fission/fetcher:latest-amd64 - fission/fetcher:latest-arm64 - - fission/fetcher:latest-armv7 - name_template: ghcr.io/fission/fission-bundle:latest image_templates: - ghcr.io/fission/fission-bundle:latest-amd64 - ghcr.io/fission/fission-bundle:latest-arm64 - - ghcr.io/fission/fission-bundle:latest-armv7 - name_template: fission/fission-bundle:latest image_templates: - fission/fission-bundle:latest-amd64 - fission/fission-bundle:latest-arm64 - - fission/fission-bundle:latest-armv7 - name_template: ghcr.io/fission/pre-upgrade-checks:latest image_templates: - ghcr.io/fission/pre-upgrade-checks:latest-amd64 - ghcr.io/fission/pre-upgrade-checks:latest-arm64 - - ghcr.io/fission/pre-upgrade-checks:latest-armv7 - name_template: fission/pre-upgrade-checks:latest image_templates: - fission/pre-upgrade-checks:latest-amd64 - fission/pre-upgrade-checks:latest-arm64 - - fission/pre-upgrade-checks:latest-armv7 - name_template: ghcr.io/fission/reporter:latest image_templates: - ghcr.io/fission/reporter:latest-amd64 - ghcr.io/fission/reporter:latest-arm64 - - ghcr.io/fission/reporter:latest-armv7 - name_template: fission/reporter:latest image_templates: - fission/reporter:latest-amd64 - fission/reporter:latest-arm64 - - fission/reporter:latest-armv7 changelog: skip: false archives: diff --git a/pkg/mqtrigger/messageQueue/kafka/consumer.go b/pkg/mqtrigger/messageQueue/kafka/consumer.go index 7cf96b9b..b684c8bc 100644 --- a/pkg/mqtrigger/messageQueue/kafka/consumer.go +++ b/pkg/mqtrigger/messageQueue/kafka/consumer.go @@ -74,6 +74,8 @@ func NewMqtConsumerGroupHandler(version sarama.KafkaVersion, // Setup implemented to satisfy the sarama.ConsumerGroupHandler interface func (ch MqtConsumerGroupHandler) Setup(session sarama.ConsumerGroupSession) error { + mqtrigger.SetTriggerStatus(ch.trigger.ObjectMeta.Name, ch.trigger.ObjectMeta.Namespace) + mqtrigger.IncreaseInprocessCount() ch.logger.With( zap.String("trigger", ch.trigger.ObjectMeta.Name), zap.String("topic", ch.trigger.Spec.Topic), @@ -88,6 +90,8 @@ func (ch MqtConsumerGroupHandler) Setup(session sarama.ConsumerGroupSession) err // Cleanup implemented to satisfy the sarama.ConsumerGroupHandler interface func (ch MqtConsumerGroupHandler) Cleanup(session sarama.ConsumerGroupSession) error { + mqtrigger.ResetTriggerStatus(ch.trigger.ObjectMeta.Name, ch.trigger.ObjectMeta.Namespace) + mqtrigger.DecreaseInprocessCount() ch.logger.With( zap.String("trigger", ch.trigger.ObjectMeta.Name), zap.String("topic", ch.trigger.Spec.Topic), diff --git a/pkg/mqtrigger/messageQueue/kafka/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go index c6ddd5e7..257aa914 100644 --- a/pkg/mqtrigger/messageQueue/kafka/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -54,7 +54,6 @@ type ( routerUrl string brokers []string version sarama.KafkaVersion - client sarama.Client authKeys map[string][]byte tls bool } @@ -111,18 +110,24 @@ func New(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messa logger.Info("created kafka queue", zap.Any("kafka brokers", kafka.brokers), zap.Any("kafka version", kafka.version)) + return kafka, nil +} - // Create new config - saramaConfig := sarama.NewConfig() - saramaConfig.Version = kafka.version +func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { + kafka.logger.Debug("inside kakfa subscribe", zap.Any("trigger", trigger)) + kafka.logger.Debug("brokers set", zap.Strings("brokers", kafka.brokers)) - // consumer config - saramaConfig.Consumer.Return.Errors = true + // Create new consumer + consumerConfig := sarama.NewConfig() + consumerConfig.Consumer.Return.Errors = true + consumerConfig.Version = kafka.version - // producer config - saramaConfig.Producer.RequiredAcks = sarama.WaitForAll - saramaConfig.Producer.Retry.Max = 10 - saramaConfig.Producer.Return.Successes = true + // Create new producer + producerConfig := sarama.NewConfig() + producerConfig.Producer.RequiredAcks = sarama.WaitForAll + producerConfig.Producer.Retry.Max = 10 + producerConfig.Producer.Return.Successes = true + producerConfig.Version = kafka.version // Setup TLS for both producer and consumer if kafka.tls { @@ -132,30 +137,18 @@ func New(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messa return nil, err } - saramaConfig.Net.TLS.Enable = true - saramaConfig.Net.TLS.Config = tlsConfig + producerConfig.Net.TLS.Enable = true + producerConfig.Net.TLS.Config = tlsConfig + consumerConfig.Net.TLS.Enable = true + consumerConfig.Net.TLS.Config = tlsConfig } - saramaClient, err := sarama.NewClient(kafka.brokers, saramaConfig) + consumer, err := sarama.NewConsumerGroup(kafka.brokers, string(trigger.ObjectMeta.UID), consumerConfig) if err != nil { return nil, err } - kafka.client = saramaClient - - return kafka, nil -} - -func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { - kafka.logger.Debug("inside kakfa subscribe", zap.Any("trigger", trigger)) - kafka.logger.Debug("brokers set", zap.Strings("brokers", kafka.brokers)) - - consumer, err := sarama.NewConsumerGroupFromClient(string(trigger.ObjectMeta.UID), kafka.client) - if err != nil { - return nil, err - } - - producer, err := sarama.NewSyncProducerFromClient(kafka.client) + producer, err := sarama.NewSyncProducer(kafka.brokers, producerConfig) if err != nil { return nil, err } diff --git a/pkg/mqtrigger/metrics.go b/pkg/mqtrigger/metrics.go index 13d14fb0..9b881ef0 100644 --- a/pkg/mqtrigger/metrics.go +++ b/pkg/mqtrigger/metrics.go @@ -45,6 +45,20 @@ var ( }, []string{"trigger_name", "trigger_namespace", "topic", "partition"}, ) + triggerStatus = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "fission_mqt_status", + Help: "Status of an individual trigger 1 if processing otherwise 0", + }, + []string{"trigger_name", "trigger_namespace"}, + ) + mqtInprocessCount = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "fission_mqt_inprocess", + Help: "Total number of MQTs in active processing", + }, + []string{}, + ) ) func IncreaseSubscriptionCount() { @@ -55,6 +69,22 @@ func DecreaseSubscriptionCount() { subscriptionCount.WithLabelValues().Dec() } +func SetTriggerStatus(trigname, trignamespace string) { + triggerStatus.WithLabelValues(trigname, trignamespace).Inc() +} + +func ResetTriggerStatus(trigname, trignamespace string) { + triggerStatus.WithLabelValues(trigname, trignamespace).Dec() +} + +func IncreaseInprocessCount() { + mqtInprocessCount.WithLabelValues().Inc() +} + +func DecreaseInprocessCount() { + mqtInprocessCount.WithLabelValues().Dec() +} + func IncreaseMessageCount(trigname, trignamespace string) { messageCount.WithLabelValues(trigname, trignamespace).Inc() } @@ -68,4 +98,6 @@ func init() { registry.MustRegister(subscriptionCount) registry.MustRegister(messageCount) registry.MustRegister(messageLagCount) + registry.MustRegister(mqtInprocessCount) + registry.MustRegister(triggerStatus) }