MQT Kafka: Use Sarama Group Consumer instead of bsm/sarama-cluster library (#2286)
Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
@@ -8,7 +8,6 @@ require (
|
|||||||
github.com/Nvveen/Gotty v0.0.0-20120604004816-cd527374f1e5 // indirect
|
github.com/Nvveen/Gotty v0.0.0-20120604004816-cd527374f1e5 // indirect
|
||||||
github.com/Shopify/sarama v1.30.0
|
github.com/Shopify/sarama v1.30.0
|
||||||
github.com/blend/go-sdk v1.20211025.3 // indirect
|
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/containerd/continuity v0.2.1 // indirect
|
||||||
github.com/dchest/uniuri v0.0.0-20200228104902-7aecb25e1fe5
|
github.com/dchest/uniuri v0.0.0-20200228104902-7aecb25e1fe5
|
||||||
github.com/docker/go-connections v0.4.0 // indirect
|
github.com/docker/go-connections v0.4.0 // indirect
|
||||||
|
|||||||
@@ -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/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/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/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/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/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=
|
github.com/casbin/casbin/v2 v2.1.2/go.mod h1:YcPU1XXisHhLzuxH9coDNf2FbKpjGlbCg3n9yuLkIJQ=
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ limitations under the License.
|
|||||||
package kafka
|
package kafka
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"crypto/tls"
|
"crypto/tls"
|
||||||
"crypto/x509"
|
"crypto/x509"
|
||||||
"fmt"
|
"fmt"
|
||||||
@@ -28,7 +29,6 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/Shopify/sarama"
|
"github.com/Shopify/sarama"
|
||||||
cluster "github.com/bsm/sarama-cluster"
|
|
||||||
"github.com/pkg/errors"
|
"github.com/pkg/errors"
|
||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
|
|
||||||
@@ -65,6 +65,183 @@ type (
|
|||||||
Factory struct{}
|
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) {
|
func (factory *Factory) Create(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) {
|
||||||
return New(logger, mqCfg, routerUrl)
|
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))
|
kafka.logger.Info("brokers set", zap.Strings("brokers", kafka.brokers))
|
||||||
|
|
||||||
// Create new consumer
|
// Create new consumer
|
||||||
consumerConfig := cluster.NewConfig()
|
consumerConfig := sarama.NewConfig()
|
||||||
consumerConfig.Consumer.Return.Errors = true
|
consumerConfig.Consumer.Return.Errors = true
|
||||||
consumerConfig.Group.Return.Notifications = true
|
consumerConfig.Version = kafka.version
|
||||||
consumerConfig.Config.Version = kafka.version
|
|
||||||
|
|
||||||
// Create new producer
|
// Create new producer
|
||||||
producerConfig := sarama.NewConfig()
|
producerConfig := sarama.NewConfig()
|
||||||
@@ -142,7 +318,8 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub
|
|||||||
consumerConfig.Net.TLS.Config = tlsConfig
|
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),
|
kafka.logger.Info("created a new consumer", zap.Strings("brokers", kafka.brokers),
|
||||||
zap.String("input topic", trigger.Spec.Topic),
|
zap.String("input topic", trigger.Spec.Topic),
|
||||||
zap.String("output topic", trigger.Spec.ResponseTopic),
|
zap.String("output topic", trigger.Spec.ResponseTopic),
|
||||||
@@ -174,18 +351,14 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub
|
|||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
// consume notifications
|
ch := NewMqtConsumerGroupHandler(kafka.version, kafka.logger, trigger, producer, kafka.routerUrl)
|
||||||
go func() {
|
|
||||||
for ntf := range consumer.Notifications() {
|
|
||||||
kafka.logger.Info("consumer notification", zap.Any("notification", ntf))
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
// consume messages
|
// consume messages
|
||||||
go func() {
|
go func() {
|
||||||
for msg := range consumer.Messages() {
|
topic := []string{trigger.Spec.Topic}
|
||||||
kafka.logger.Debug("calling message handler", zap.String("message", string(msg.Value[:])))
|
ctx := context.Background()
|
||||||
go kafkaMsgHandler(&kafka, producer, trigger, msg, consumer)
|
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 {
|
func (kafka Kafka) Unsubscribe(subscription messageQueue.Subscription) error {
|
||||||
return subscription.(*cluster.Consumer).Close()
|
return subscription.(sarama.ConsumerGroup).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
|
|
||||||
}
|
}
|
||||||
|
|
||||||
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) {
|
||||||
|
|||||||
Reference in New Issue
Block a user