From 2732aff425f68aeb25796d383179ce41b35a4dc0 Mon Sep 17 00:00:00 2001 From: Rutul Patel Date: Wed, 9 Mar 2022 15:47:05 +0530 Subject: [PATCH 01/13] Added pprof profiling to monitor heap memory --- cmd/service/metro/metro.go | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) diff --git a/cmd/service/metro/metro.go b/cmd/service/metro/metro.go index 94cc570c..33aa4054 100644 --- a/cmd/service/metro/metro.go +++ b/cmd/service/metro/metro.go @@ -12,6 +12,9 @@ import ( configreader "github.com/razorpay/metro/pkg/config" "github.com/razorpay/metro/pkg/encryption" "github.com/razorpay/metro/pkg/logger" + + "net/http" + _ "net/http/pprof" ) const ( @@ -38,7 +41,7 @@ func isValidComponent(component string) bool { } // Init initializes all modules (logger, tracing, config, metro component) -func Init(_ context.Context, env string, componentName string) { +func Init(ctx context.Context, env string, componentName string) { // componentName validation ok := isValidComponent(componentName) if !ok { @@ -68,6 +71,8 @@ func Init(_ context.Context, env string, componentName string) { err = boot.InitMonitoring(env, appConfig.App, appConfig.Sentry, appConfig.Tracing) + setPprofProfiles(ctx, componentName) + if err != nil { log.Fatalf("error in setting up monitoring : %v", err) } @@ -114,3 +119,15 @@ func Run(ctx context.Context) { logger.Ctx(ctx).Infow("stopped metro") } + +// sets up pprof profile for perfomance monitoring +func setPprofProfiles(ctx context.Context, componentName string) { + logger.Ctx(ctx).Infow("initialising pprof profiles") + go func() { + if componentName == Web { + http.ListenAndServe("metro-web-pprof.concierge.stage.razorpay.in:8080", nil) + } else if componentName == Worker { + http.ListenAndServe("metro-worker-pprof.concierge.stage.razorpay.in:8080", nil) + } + }() +} From 890412adc765bc572f64b801c60600351cee1108 Mon Sep 17 00:00:00 2001 From: Rutul Patel Date: Wed, 9 Mar 2022 15:56:11 +0530 Subject: [PATCH 02/13] lint check improvement --- cmd/service/metro/metro.go | 1 + 1 file changed, 1 insertion(+) diff --git a/cmd/service/metro/metro.go b/cmd/service/metro/metro.go index 33aa4054..c35bfab3 100644 --- a/cmd/service/metro/metro.go +++ b/cmd/service/metro/metro.go @@ -14,6 +14,7 @@ import ( "github.com/razorpay/metro/pkg/logger" "net/http" + // blank import added for testing. _ "net/http/pprof" ) From ff610644e2a65285faf0d57aea516f4fd4439d4d Mon Sep 17 00:00:00 2001 From: Rutul Patel Date: Wed, 9 Mar 2022 16:10:39 +0530 Subject: [PATCH 03/13] added port for pprof --- build/docker/Dockerfile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/build/docker/Dockerfile b/build/docker/Dockerfile index 5e0a73b8..2e9abac7 100644 --- a/build/docker/Dockerfile +++ b/build/docker/Dockerfile @@ -78,7 +78,7 @@ WORKDIR /app RUN apk add --update --no-cache dumb-init su-exec ca-certificates curl RUN mkdir -p /app/public -EXPOSE 8081 8082 8083 8084 3000 +EXPOSE 8081 8082 8083 8084 3000 8080 RUN chmod +x entrypoint.sh probe.sh ENTRYPOINT ["/app/entrypoint.sh", "metro"] From 3cdd82f4f96cfb4f033ab0e29f8c0ad10eb88ccf Mon Sep 17 00:00:00 2001 From: Rutul Patel Date: Wed, 9 Mar 2022 16:53:38 +0530 Subject: [PATCH 04/13] wip --- cmd/service/metro/metro.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/cmd/service/metro/metro.go b/cmd/service/metro/metro.go index c35bfab3..c79c6462 100644 --- a/cmd/service/metro/metro.go +++ b/cmd/service/metro/metro.go @@ -14,6 +14,7 @@ import ( "github.com/razorpay/metro/pkg/logger" "net/http" + // blank import added for testing. _ "net/http/pprof" ) @@ -125,10 +126,9 @@ func Run(ctx context.Context) { func setPprofProfiles(ctx context.Context, componentName string) { logger.Ctx(ctx).Infow("initialising pprof profiles") go func() { - if componentName == Web { - http.ListenAndServe("metro-web-pprof.concierge.stage.razorpay.in:8080", nil) - } else if componentName == Worker { - http.ListenAndServe("metro-worker-pprof.concierge.stage.razorpay.in:8080", nil) + myMux := http.DefaultServeMux + if err := http.ListenAndServe("localhost:8080", myMux); err != nil { + logger.Ctx(ctx).Fatalw("Error when starting or running %v pprof http server: %v", componentName, err) } }() } From 0f104801853fe0e5c6d658f807759573efa134e8 Mon Sep 17 00:00:00 2001 From: vnktram Date: Fri, 25 Mar 2022 13:06:51 +0530 Subject: [PATCH 05/13] kafka client configs --- pkg/messagebroker/kafka.go | 21 ++++++++++++--------- 1 file changed, 12 insertions(+), 9 deletions(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index 7f4fd4c7..16a82c20 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -109,15 +109,18 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options logger.Ctx(ctx).Infow("kafka producer: initializing new", "options", options) configMap := &kafkapkg.ConfigMap{ - "bootstrap.servers": strings.Join(bConfig.Brokers, ","), - "socket.keepalive.enable": true, - "retries": 3, - "linger.ms": 0, - "request.timeout.ms": 3000, - "delivery.timeout.ms": 10000, - "connections.max.idle.ms": 180000, - "go.logs.channel.enable": true, - "debug": "all", + "bootstrap.servers": strings.Join(bConfig.Brokers, ","), + "socket.keepalive.enable": true, + "acks": 1, + "retries": 3, + "linger.ms": 0, + "request.timeout.ms": 3000, + "delivery.timeout.ms": 10000, + "connections.max.idle.ms": 180000, + "log.queue": false, + "queue.buffering.max.messages": 100, + "go.logs.channel.enable": true, + "debug": "all", } if bConfig.EnableTLS { From faecdaf813fcdfb2dc3f7e581bc67fe9a1dcdcb5 Mon Sep 17 00:00:00 2001 From: vnktram Date: Tue, 29 Mar 2022 14:46:28 +0530 Subject: [PATCH 06/13] Introduce batch and report configs --- pkg/messagebroker/kafka.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index 16a82c20..dc2c660d 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -119,7 +119,9 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "connections.max.idle.ms": 180000, "log.queue": false, "queue.buffering.max.messages": 100, - "go.logs.channel.enable": true, + "go.logs.channel.enable": false, + "go.delivery.reports": false, + "go.batch.producer": true, "debug": "all", } From ec318d7fe1ac6612ba15bdb195464bdb1b919a06 Mon Sep 17 00:00:00 2001 From: vnktram Date: Mon, 4 Apr 2022 19:42:13 +0530 Subject: [PATCH 07/13] Restrict channel size --- pkg/messagebroker/kafka.go | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index dc2c660d..ffc4602d 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -120,6 +120,8 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "log.queue": false, "queue.buffering.max.messages": 100, "go.logs.channel.enable": false, + "go.events.channel.size": 1, + "go.produce.channel.size": 1000, "go.delivery.reports": false, "go.batch.producer": true, "debug": "all", @@ -383,7 +385,7 @@ func (k *KafkaBroker) SendMessage(ctx context.Context, request SendMessageToTopi logger.Ctx(ctx).Warnw("error injecting span context in message headers", "error", injectErr.Error()) } - deliveryChan := make(chan kafkapkg.Event, 1000) + deliveryChan := make(chan kafkapkg.Event, 1) defer close(deliveryChan) topicN := normalizeTopicName(request.Topic) @@ -405,10 +407,9 @@ func (k *KafkaBroker) SendMessage(ctx context.Context, request SendMessageToTopi } var m *kafkapkg.Message - select { - case event := <-deliveryChan: - m = event.(*kafkapkg.Message) - } + + event := <-deliveryChan + m = event.(*kafkapkg.Message) if m != nil && m.TopicPartition.Error != nil { logger.Ctx(ctx).Errorw("kafka: error in publishing messages", "error", m.TopicPartition.Error.Error()) From 58bdf45b888fd043af27514f0690e3c8a5804370 Mon Sep 17 00:00:00 2001 From: vnktram Date: Mon, 4 Apr 2022 19:43:18 +0530 Subject: [PATCH 08/13] Enable batch producer --- pkg/messagebroker/kafka.go | 1 + 1 file changed, 1 insertion(+) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index ffc4602d..d7f07b66 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -113,6 +113,7 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "socket.keepalive.enable": true, "acks": 1, "retries": 3, + "go.batch.producer": true, "linger.ms": 0, "request.timeout.ms": 3000, "delivery.timeout.ms": 10000, From b3385f64457ced090430a4711e1f39f61c06187b Mon Sep 17 00:00:00 2001 From: vnktram Date: Mon, 4 Apr 2022 19:43:43 +0530 Subject: [PATCH 09/13] Remove duplicate key --- pkg/messagebroker/kafka.go | 1 - 1 file changed, 1 deletion(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index d7f07b66..00ad352a 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -124,7 +124,6 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "go.events.channel.size": 1, "go.produce.channel.size": 1000, "go.delivery.reports": false, - "go.batch.producer": true, "debug": "all", } From b7efb6a61de92f1d005f14ccc25340a82e764300 Mon Sep 17 00:00:00 2001 From: vnktram Date: Tue, 5 Apr 2022 22:54:39 +0530 Subject: [PATCH 10/13] Remove batch producer mode --- pkg/messagebroker/kafka.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index 00ad352a..870654f5 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -113,7 +113,7 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "socket.keepalive.enable": true, "acks": 1, "retries": 3, - "go.batch.producer": true, + "go.batch.producer": false, "linger.ms": 0, "request.timeout.ms": 3000, "delivery.timeout.ms": 10000, From 9571cc1b842f46cc0f771477b540f68575c2fa8d Mon Sep 17 00:00:00 2001 From: vnktram Date: Tue, 5 Apr 2022 23:26:05 +0530 Subject: [PATCH 11/13] Change debug level --- pkg/messagebroker/kafka.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index 870654f5..c9caa886 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -124,7 +124,7 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "go.events.channel.size": 1, "go.produce.channel.size": 1000, "go.delivery.reports": false, - "debug": "all", + "debug": "error", } if bConfig.EnableTLS { From 9a7a49e9c603ba6c47a3510d3be2caf3c239b24e Mon Sep 17 00:00:00 2001 From: vnktram Date: Tue, 5 Apr 2022 23:34:49 +0530 Subject: [PATCH 12/13] Remove debug logs --- pkg/messagebroker/kafka.go | 1 - 1 file changed, 1 deletion(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index c9caa886..30178eb4 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -124,7 +124,6 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "go.events.channel.size": 1, "go.produce.channel.size": 1000, "go.delivery.reports": false, - "debug": "error", } if bConfig.EnableTLS { From 41dddd8010a356953792869a6649f196095c805a Mon Sep 17 00:00:00 2001 From: vnktram Date: Wed, 6 Apr 2022 00:13:02 +0530 Subject: [PATCH 13/13] Increase channel size --- pkg/messagebroker/kafka.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg/messagebroker/kafka.go b/pkg/messagebroker/kafka.go index 30178eb4..af0b357a 100644 --- a/pkg/messagebroker/kafka.go +++ b/pkg/messagebroker/kafka.go @@ -119,9 +119,9 @@ func newKafkaProducerClient(ctx context.Context, bConfig *BrokerConfig, options "delivery.timeout.ms": 10000, "connections.max.idle.ms": 180000, "log.queue": false, - "queue.buffering.max.messages": 100, + "queue.buffering.max.messages": 1000, "go.logs.channel.enable": false, - "go.events.channel.size": 1, + "go.events.channel.size": 100, "go.produce.channel.size": 1000, "go.delivery.reports": false, }