Optimize Kafka Client in Kafka Connector (#2630)
Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
@@ -106,8 +106,8 @@ func (ch MqtConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSessi
|
||||
topic := claim.Topic()
|
||||
partition := string(claim.Partition())
|
||||
|
||||
// initially set metrics to -1
|
||||
mqtrigger.SetMessageLagCount(trigger, triggerNamespace, topic, partition, -1)
|
||||
// initially set message lag count
|
||||
mqtrigger.SetMessageLagCount(trigger, triggerNamespace, topic, partition, claim.HighWaterMarkOffset()-claim.InitialOffset())
|
||||
|
||||
// Do not move the code below to a goroutine.
|
||||
// The `ConsumeClaim` itself is called within a goroutine
|
||||
|
||||
@@ -54,6 +54,7 @@ type (
|
||||
routerUrl string
|
||||
brokers []string
|
||||
version sarama.KafkaVersion
|
||||
client sarama.Client
|
||||
authKeys map[string][]byte
|
||||
tls bool
|
||||
}
|
||||
@@ -110,24 +111,18 @@ 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
|
||||
}
|
||||
|
||||
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))
|
||||
// Create new config
|
||||
saramaConfig := sarama.NewConfig()
|
||||
saramaConfig.Version = kafka.version
|
||||
|
||||
// Create new consumer
|
||||
consumerConfig := sarama.NewConfig()
|
||||
consumerConfig.Consumer.Return.Errors = true
|
||||
consumerConfig.Version = kafka.version
|
||||
// consumer config
|
||||
saramaConfig.Consumer.Return.Errors = 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
|
||||
// producer config
|
||||
saramaConfig.Producer.RequiredAcks = sarama.WaitForAll
|
||||
saramaConfig.Producer.Retry.Max = 10
|
||||
saramaConfig.Producer.Return.Successes = true
|
||||
|
||||
// Setup TLS for both producer and consumer
|
||||
if kafka.tls {
|
||||
@@ -137,18 +132,30 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub
|
||||
return nil, err
|
||||
}
|
||||
|
||||
producerConfig.Net.TLS.Enable = true
|
||||
producerConfig.Net.TLS.Config = tlsConfig
|
||||
consumerConfig.Net.TLS.Enable = true
|
||||
consumerConfig.Net.TLS.Config = tlsConfig
|
||||
saramaConfig.Net.TLS.Enable = true
|
||||
saramaConfig.Net.TLS.Config = tlsConfig
|
||||
}
|
||||
|
||||
consumer, err := sarama.NewConsumerGroup(kafka.brokers, string(trigger.ObjectMeta.UID), consumerConfig)
|
||||
saramaClient, err := sarama.NewClient(kafka.brokers, saramaConfig)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
producer, err := sarama.NewSyncProducer(kafka.brokers, producerConfig)
|
||||
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)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user