From bf8c01cff0265f0090c9d81c734e6164436de08b Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Wed, 22 Dec 2021 18:00:25 +0530 Subject: [PATCH] MQT Kafka: Use Sarama Group Consumer instead of bsm/sarama-cluster library (#2286) Signed-off-by: Sanket Sudake --- go.mod | 1 - go.sum | 2 - pkg/mqtrigger/messageQueue/kafka/kafka.go | 344 ++++++++++++---------- 3 files changed, 189 insertions(+), 158 deletions(-) diff --git a/go.mod b/go.mod index 6aa9d1ad..e0bdcb8a 100644 --- a/go.mod +++ b/go.mod @@ -8,7 +8,6 @@ require ( github.com/Nvveen/Gotty v0.0.0-20120604004816-cd527374f1e5 // indirect github.com/Shopify/sarama v1.30.0 github.com/blend/go-sdk v1.20211025.3 // indirect - github.com/bsm/sarama-cluster v2.1.15+incompatible github.com/containerd/continuity v0.2.1 // indirect github.com/dchest/uniuri v0.0.0-20200228104902-7aecb25e1fe5 github.com/docker/go-connections v0.4.0 // indirect diff --git a/go.sum b/go.sum index 73bed072..c7a91951 100644 --- a/go.sum +++ b/go.sum @@ -187,8 +187,6 @@ github.com/blend/sentry-go v1.0.1/go.mod h1:hgyX3WXen2YBiA0NitlfsXsvS+9ly2YlEBmm github.com/bmizerany/pat v0.0.0-20170815010413-6226ea591a40/go.mod h1:8rLXio+WjiTceGBHIoTvn60HIbs7Hm7bcHjyrSqYB9c= github.com/boltdb/bolt v1.3.1/go.mod h1:clJnj/oiGkjum5o1McbSZDSLxVThjynRyGBgiAx27Ps= github.com/bonitoo-io/go-sql-bigquery v0.3.4-1.4.0/go.mod h1:J4Y6YJm0qTWB9aFziB7cPeSyc6dOZFyJdteSeybVpXQ= -github.com/bsm/sarama-cluster v2.1.15+incompatible h1:RkV6WiNRnqEEbp81druK8zYhmnIgdOjqSVi0+9Cnl2A= -github.com/bsm/sarama-cluster v2.1.15+incompatible/go.mod h1:r7ao+4tTNXvWm+VRpRJchr2kQhqxgmAp2iEX5W96gMM= github.com/c-bata/go-prompt v0.2.2/go.mod h1:VzqtzE2ksDBcdln8G7mk2RX9QyGjH+OVqOCSiVIqS34= github.com/cactus/go-statsd-client/statsd v0.0.0-20191106001114-12b4e2b38748/go.mod h1:l/bIBLeOl9eX+wxJAzxS4TveKRtAqlyDpHjhkfO0MEI= github.com/casbin/casbin/v2 v2.1.2/go.mod h1:YcPU1XXisHhLzuxH9coDNf2FbKpjGlbCg3n9yuLkIJQ= diff --git a/pkg/mqtrigger/messageQueue/kafka/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go index 466cd5d1..0ced9ac5 100644 --- a/pkg/mqtrigger/messageQueue/kafka/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -17,6 +17,7 @@ limitations under the License. package kafka import ( + "context" "crypto/tls" "crypto/x509" "fmt" @@ -28,7 +29,6 @@ import ( "strings" "github.com/Shopify/sarama" - cluster "github.com/bsm/sarama-cluster" "github.com/pkg/errors" "go.uber.org/zap" @@ -65,6 +65,183 @@ type ( Factory struct{} ) +type MqtConsumerGroupHandler struct { + version sarama.KafkaVersion + logger *zap.Logger + trigger *fv1.MessageQueueTrigger + fissionHeaders map[string]string + producer sarama.SyncProducer + fnUrl string +} + +func NewMqtConsumerGroupHandler(version sarama.KafkaVersion, + logger *zap.Logger, + trigger *fv1.MessageQueueTrigger, + producer sarama.SyncProducer, + routerUrl string) MqtConsumerGroupHandler { + ch := MqtConsumerGroupHandler{ + version: version, + logger: logger, + trigger: trigger, + producer: producer, + } + // Support other function ref types + if ch.trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { + ch.logger.Fatal("unsupported function reference type for trigger", + zap.Any("function_reference_type", ch.trigger.Spec.FunctionReference.Type), + zap.String("trigger", ch.trigger.ObjectMeta.Name)) + } + // Generate the Headers + ch.fissionHeaders = map[string]string{ + "X-Fission-MQTrigger-Topic": ch.trigger.Spec.Topic, + "X-Fission-MQTrigger-RespTopic": ch.trigger.Spec.ResponseTopic, + "X-Fission-MQTrigger-ErrorTopic": ch.trigger.Spec.ErrorTopic, + "Content-Type": ch.trigger.Spec.ContentType, + } + ch.fnUrl = routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(ch.trigger.Spec.FunctionReference.Name, ch.trigger.ObjectMeta.Namespace), "/") + ch.logger.Debug("function HTTP URL", zap.String("url", ch.fnUrl)) + return ch +} + +func (ch MqtConsumerGroupHandler) Setup(sarama.ConsumerGroupSession) error { + return nil +} + +func (ch MqtConsumerGroupHandler) Cleanup(sarama.ConsumerGroupSession) error { + return nil +} + +func (ch MqtConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { + for msg := range claim.Messages() { + ch.kafkaMsgHandler(session, msg) + } + return nil +} + +//func (ch *MqtConsumerGroupHandler) kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.MessageQueueTrigger, msg *sarama.ConsumerMessage, consumer *cluster.Consumer) { + +func (ch *MqtConsumerGroupHandler) kafkaMsgHandler(session sarama.ConsumerGroupSession, msg *sarama.ConsumerMessage) { + var value string = string(msg.Value[:]) + + // Create request + req, err := http.NewRequest("POST", ch.fnUrl, strings.NewReader(value)) + if err != nil { + ch.logger.Error("failed to create HTTP request to invoke function", + zap.Error(err), + zap.String("function_url", ch.fnUrl)) + return + } + + // Set the headers came from Kafka record + // Using Header.Add() as msg.Headers may have keys with more than one value + if ch.version.IsAtLeast(sarama.V0_11_0_0) { + for _, h := range msg.Headers { + req.Header.Add(string(h.Key), string(h.Value)) + } + } else { + ch.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", + zap.Any("current_version", ch.version)) + } + + for k, v := range ch.fissionHeaders { + req.Header.Set(k, v) + } + + // Make the request + var resp *http.Response + for attempt := 0; attempt <= ch.trigger.Spec.MaxRetries; attempt++ { + // Make the request + resp, err = http.DefaultClient.Do(req) + if err != nil { + ch.logger.Error("sending function invocation request failed", + zap.Error(err), + zap.String("function_url", ch.fnUrl), + zap.String("trigger", ch.trigger.ObjectMeta.Name)) + continue + } + if resp == nil { + continue + } + if err == nil && resp.StatusCode == http.StatusOK { + // Success, quit retrying + break + } + } + + generateErrorHeaders := func(errString string) []sarama.RecordHeader { + var errorHeaders []sarama.RecordHeader + if ch.version.IsAtLeast(sarama.V0_11_0_0) { + if count, ok := errorMessageMap[errString]; ok { + errorMessageMap[errString] = count + 1 + } else { + errorMessageMap[errString] = 1 + } + errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("MessageSource"), Value: []byte(ch.trigger.Spec.Topic)}) + errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("RecycleCounter"), Value: []byte(strconv.Itoa(errorMessageMap[errString]))}) + } + return errorHeaders + } + + if resp == nil { + errorString := fmt.Sprintf("request exceed retries: %v", ch.trigger.Spec.MaxRetries) + errorHeaders := generateErrorHeaders(errorString) + errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, + fmt.Errorf(errorString), errorHeaders) + return + } + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + + ch.logger.Debug("got response from function invocation", + zap.String("function_url", ch.fnUrl), + zap.String("trigger", ch.trigger.ObjectMeta.Name), + zap.String("body", string(body))) + + if err != nil { + errorString := "request body error: " + string(body) + errorHeaders := generateErrorHeaders(errorString) + errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, + errors.Wrapf(err, errorString), errorHeaders) + return + } + if resp.StatusCode != 200 { + errorString := fmt.Sprintf("request returned failure: %v, request body error: %v", resp.StatusCode, body) + errorHeaders := generateErrorHeaders(errorString) + errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, + fmt.Errorf("request returned failure: %v", resp.StatusCode), errorHeaders) + return + } + if len(ch.trigger.Spec.ResponseTopic) > 0 { + // Generate Kafka record headers + var kafkaRecordHeaders []sarama.RecordHeader + if ch.version.IsAtLeast(sarama.V0_11_0_0) { + for k, v := range resp.Header { + // One key may have multiple values + for _, v := range v { + kafkaRecordHeaders = append(kafkaRecordHeaders, sarama.RecordHeader{Key: []byte(k), Value: []byte(v)}) + } + } + } else { + ch.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", + zap.Any("current_version", ch.version)) + } + + _, _, err := ch.producer.SendMessage(&sarama.ProducerMessage{ + Topic: ch.trigger.Spec.ResponseTopic, + Value: sarama.StringEncoder(body), + Headers: kafkaRecordHeaders, + }) + if err != nil { + ch.logger.Warn("failed to publish response body from function invocation to topic", + zap.Error(err), + zap.String("topic", ch.trigger.Spec.Topic), + zap.String("function_url", ch.fnUrl)) + return + } + } + session.MarkMessage(msg, "") +} + func (factory *Factory) Create(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { return New(logger, mqCfg, routerUrl) } @@ -116,10 +293,9 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub kafka.logger.Info("brokers set", zap.Strings("brokers", kafka.brokers)) // Create new consumer - consumerConfig := cluster.NewConfig() + consumerConfig := sarama.NewConfig() consumerConfig.Consumer.Return.Errors = true - consumerConfig.Group.Return.Notifications = true - consumerConfig.Config.Version = kafka.version + consumerConfig.Version = kafka.version // Create new producer producerConfig := sarama.NewConfig() @@ -142,7 +318,8 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub consumerConfig.Net.TLS.Config = tlsConfig } - consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.ObjectMeta.UID), []string{trigger.Spec.Topic}, consumerConfig) + consumer, err := sarama.NewConsumerGroup(kafka.brokers, string(trigger.ObjectMeta.UID), consumerConfig) + // consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.ObjectMeta.UID), []string{trigger.Spec.Topic}, consumerConfig) kafka.logger.Info("created a new consumer", zap.Strings("brokers", kafka.brokers), zap.String("input topic", trigger.Spec.Topic), zap.String("output topic", trigger.Spec.ResponseTopic), @@ -174,18 +351,14 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub } }() - // consume notifications - go func() { - for ntf := range consumer.Notifications() { - kafka.logger.Info("consumer notification", zap.Any("notification", ntf)) - } - }() - + ch := NewMqtConsumerGroupHandler(kafka.version, kafka.logger, trigger, producer, kafka.routerUrl) // consume messages go func() { - for msg := range consumer.Messages() { - kafka.logger.Debug("calling message handler", zap.String("message", string(msg.Value[:]))) - go kafkaMsgHandler(&kafka, producer, trigger, msg, consumer) + 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)) } }() @@ -217,146 +390,7 @@ func (kafka Kafka) getTLSConfig() (*tls.Config, error) { } func (kafka Kafka) Unsubscribe(subscription messageQueue.Subscription) error { - return subscription.(*cluster.Consumer).Close() -} - -func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.MessageQueueTrigger, msg *sarama.ConsumerMessage, consumer *cluster.Consumer) { - var value string = string(msg.Value[:]) - // Support other function ref types - if trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { - kafka.logger.Fatal("unsupported function reference type for trigger", - zap.Any("function_reference_type", trigger.Spec.FunctionReference.Type), - zap.String("trigger", trigger.ObjectMeta.Name)) - } - - url := kafka.routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.ObjectMeta.Namespace), "/") - kafka.logger.Debug("making HTTP request", zap.String("url", url)) - - // Generate the Headers - fissionHeaders := map[string]string{ - "X-Fission-MQTrigger-Topic": trigger.Spec.Topic, - "X-Fission-MQTrigger-RespTopic": trigger.Spec.ResponseTopic, - "X-Fission-MQTrigger-ErrorTopic": trigger.Spec.ErrorTopic, - "Content-Type": trigger.Spec.ContentType, - } - - // Create request - req, err := http.NewRequest("POST", url, strings.NewReader(value)) - if err != nil { - kafka.logger.Error("failed to create HTTP request to invoke function", - zap.Error(err), - zap.String("function_url", url)) - return - } - - // Set the headers came from Kafka record - // Using Header.Add() as msg.Headers may have keys with more than one value - if kafka.version.IsAtLeast(sarama.V0_11_0_0) { - for _, h := range msg.Headers { - req.Header.Add(string(h.Key), string(h.Value)) - } - } else { - kafka.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", - zap.Any("current_version", kafka.version)) - } - - for k, v := range fissionHeaders { - req.Header.Set(k, v) - } - - // Make the request - var resp *http.Response - for attempt := 0; attempt <= trigger.Spec.MaxRetries; attempt++ { - // Make the request - resp, err = http.DefaultClient.Do(req) - if err != nil { - kafka.logger.Error("sending function invocation request failed", - zap.Error(err), - zap.String("function_url", url), - zap.String("trigger", trigger.ObjectMeta.Name)) - continue - } - if resp == nil { - continue - } - if err == nil && resp.StatusCode == http.StatusOK { - // Success, quit retrying - break - } - } - - generateErrorHeaders := func(errString string) []sarama.RecordHeader { - var errorHeaders []sarama.RecordHeader - if kafka.version.IsAtLeast(sarama.V0_11_0_0) { - if count, ok := errorMessageMap[errString]; ok { - errorMessageMap[errString] = count + 1 - } else { - errorMessageMap[errString] = 1 - } - errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("MessageSource"), Value: []byte(trigger.Spec.Topic)}) - errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("RecycleCounter"), Value: []byte(strconv.Itoa(errorMessageMap[errString]))}) - } - return errorHeaders - } - - if resp == nil { - errorString := fmt.Sprintf("request exceed retries: %v", trigger.Spec.MaxRetries) - errorHeaders := generateErrorHeaders(errorString) - errorHandler(kafka.logger, trigger, producer, url, - fmt.Errorf(errorString), errorHeaders) - return - } - defer resp.Body.Close() - body, err := io.ReadAll(resp.Body) - - kafka.logger.Debug("got response from function invocation", - zap.String("function_url", url), - zap.String("trigger", trigger.ObjectMeta.Name), - zap.String("body", string(body))) - - if err != nil { - errorString := "request body error: " + string(body) - errorHeaders := generateErrorHeaders(errorString) - errorHandler(kafka.logger, trigger, producer, url, - errors.Wrapf(err, errorString), errorHeaders) - return - } - if resp.StatusCode != 200 { - errorString := fmt.Sprintf("request returned failure: %v, request body error: %v", resp.StatusCode, body) - errorHeaders := generateErrorHeaders(errorString) - errorHandler(kafka.logger, trigger, producer, url, - fmt.Errorf("request returned failure: %v", resp.StatusCode), errorHeaders) - return - } - if len(trigger.Spec.ResponseTopic) > 0 { - // Generate Kafka record headers - var kafkaRecordHeaders []sarama.RecordHeader - if kafka.version.IsAtLeast(sarama.V0_11_0_0) { - for k, v := range resp.Header { - // One key may have multiple values - for _, v := range v { - kafkaRecordHeaders = append(kafkaRecordHeaders, sarama.RecordHeader{Key: []byte(k), Value: []byte(v)}) - } - } - } else { - kafka.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", - zap.Any("current_version", kafka.version)) - } - - _, _, err := producer.SendMessage(&sarama.ProducerMessage{ - Topic: trigger.Spec.ResponseTopic, - Value: sarama.StringEncoder(body), - Headers: kafkaRecordHeaders, - }) - if err != nil { - kafka.logger.Warn("failed to publish response body from function invocation to topic", - zap.Error(err), - zap.String("topic", trigger.Spec.Topic), - zap.String("function_url", url)) - return - } - } - consumer.MarkOffset(msg, "") // mark message as processed + return subscription.(sarama.ConsumerGroup).Close() } func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer sarama.SyncProducer, funcUrl string, err error, errorTopicHeaders []sarama.RecordHeader) {