diff --git a/extensions/kafka_producer.go b/extensions/kafka_producer.go index 98b941f..037d867 100644 --- a/extensions/kafka_producer.go +++ b/extensions/kafka_producer.go @@ -24,16 +24,36 @@ package extensions import ( "strings" + "sync" "time" - "github.com/DataDog/datadog-go/statsd" "github.com/Shopify/sarama" + "github.com/prometheus/client_golang/prometheus" "github.com/spf13/viper" "github.com/topfreegames/marathon/log" "github.com/topfreegames/marathon/messages" "github.com/uber-go/zap" ) +var ( + registerKafkaMetricsOnce sync.Once + + kafkaSendMessageReturn = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "marathon_kafka_send_message_return_total", + Help: "Kafka async producer message acknowledgements, labelled by error status.", + }, + []string{"error"}, + ) +) + +// MustRegisterKafkaMetrics registers Kafka producer Prometheus collectors. Safe to call more than once. +func MustRegisterKafkaMetrics() { + registerKafkaMetricsOnce.Do(func() { + prometheus.MustRegister(kafkaSendMessageReturn) + }) +} + // KafkaProducer is the struct that connects to Kafka type KafkaProducer struct { Config *viper.Viper @@ -43,20 +63,18 @@ type KafkaProducer struct { FlushMaxMessages int FlushFrequency int // ms Producer sarama.AsyncProducer - Statsd *statsd.Client MaxMessageBytes int Retries int } // NewKafkaProducer creates a new kafka producer -func NewKafkaProducer(config *viper.Viper, logger zap.Logger, statsd *statsd.Client) (*KafkaProducer, error) { +func NewKafkaProducer(config *viper.Viper, logger zap.Logger) (*KafkaProducer, error) { l := logger.With( zap.String("source", "KafkaExtension"), ) client := &KafkaProducer{ Config: config, Logger: l, - Statsd: statsd, } client.loadConfigurationDefaults() @@ -105,13 +123,13 @@ func (c *KafkaProducer) connectToKafka() error { go func() { for range producer.Successes() { - c.Statsd.Incr("send_message_return", []string{"error:false"}, 1) + kafkaSendMessageReturn.WithLabelValues("false").Inc() } }() go func() { for range producer.Errors() { - c.Statsd.Incr("send_message_return", []string{"error:true"}, 1) + kafkaSendMessageReturn.WithLabelValues("true").Inc() } }() diff --git a/extensions/kafka_producer_test.go b/extensions/kafka_producer_test.go index 6d810e4..d147c52 100644 --- a/extensions/kafka_producer_test.go +++ b/extensions/kafka_producer_test.go @@ -25,7 +25,6 @@ import ( "encoding/json" "time" - "github.com/DataDog/datadog-go/statsd" "github.com/confluentinc/confluent-kafka-go/kafka" . "github.com/onsi/ginkgo" . "github.com/onsi/gomega" @@ -77,7 +76,6 @@ var _ = XDescribe("Kafka Extension", func() { var logger zap.Logger var config *viper.Viper var testConsumer *kafka.Consumer - var statsdClient *statsd.Client BeforeEach(func() { logger = zap.New( @@ -96,9 +94,6 @@ var _ = XDescribe("Kafka Extension", func() { Expect(err).NotTo(HaveOccurred()) err = waitForConsumer(testConsumer) Expect(err).NotTo(HaveOccurred()) - - statsdClient, err = statsd.New("localhost:1234") - Expect(err).NotTo(HaveOccurred()) }) AfterEach(func() { @@ -108,7 +103,7 @@ var _ = XDescribe("Kafka Extension", func() { Describe("Creating new client", func() { It("should return connected client", func() { - kafka, err := extensions.NewKafkaProducer(config, logger, statsdClient) + kafka, err := extensions.NewKafkaProducer(config, logger) Expect(err).NotTo(HaveOccurred()) defer kafka.Close() @@ -119,7 +114,7 @@ var _ = XDescribe("Kafka Extension", func() { Describe("Send GCM Message", func() { It("should send GCM message", func() { - kafka, err := extensions.NewKafkaProducer(config, logger, statsdClient) + kafka, err := extensions.NewKafkaProducer(config, logger) Expect(err).NotTo(HaveOccurred()) defer kafka.Close() @@ -143,7 +138,7 @@ var _ = XDescribe("Kafka Extension", func() { Describe("Send APNS Message", func() { It("should send APNS message", func() { - kafka, err := extensions.NewKafkaProducer(config, logger, statsdClient) + kafka, err := extensions.NewKafkaProducer(config, logger) Expect(err).NotTo(HaveOccurred()) defer kafka.Close() diff --git a/go.mod b/go.mod index e036bf9..5685688 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,6 @@ module github.com/topfreegames/marathon go 1.19 require ( - github.com/DataDog/datadog-go v0.0.0-20180330214955-e67964b4021a github.com/Shopify/sarama v1.22.1 github.com/asaskevich/govalidator v0.0.0-20161001163130-7b3beb6df3c4 github.com/aws/aws-sdk-go v1.12.72 @@ -20,6 +19,8 @@ require ( github.com/onsi/ginkgo v1.14.2 github.com/onsi/gomega v1.10.4 github.com/pressly/goose v0.0.0-20161106184528-d6e8fe029271 + github.com/prometheus/client_golang v1.14.0 + github.com/prometheus/client_model v0.3.0 github.com/satori/go.uuid v1.2.0 github.com/sendgrid/sendgrid-go v3.4.1+incompatible github.com/sirupsen/logrus v1.6.0 @@ -72,8 +73,6 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pelletier/go-toml v1.2.0 // indirect github.com/pierrec/lz4 v0.0.0-20190327172049-315a67e90e41 // indirect - github.com/prometheus/client_golang v1.14.0 // indirect - github.com/prometheus/client_model v0.3.0 // indirect github.com/prometheus/common v0.37.0 // indirect github.com/prometheus/procfs v0.8.0 // indirect github.com/rcrowley/go-metrics v0.0.0-20181016184325-3113b8401b8a // indirect diff --git a/go.sum b/go.sum index 46e7971..b8a37e4 100644 --- a/go.sum +++ b/go.sum @@ -36,8 +36,6 @@ github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03 github.com/BurntSushi/toml v1.2.1 h1:9F2/+DoOYIOksmaJFPw1tGFy1eDnIJXg+UHjuD8lTak= github.com/BurntSushi/toml v1.2.1/go.mod h1:CxXYINrC8qIiEnFrOxCa7Jy5BFHlXnUU2pbicEuybxQ= github.com/BurntSushi/xgb v0.0.0-20160522181843-27f122750802/go.mod h1:IVnqGOEym/WlBOVXweHU+Q+/VP0lqqI8lqeDx9IjBqo= -github.com/DataDog/datadog-go v0.0.0-20180330214955-e67964b4021a h1:zpQSzEApXM0qkXcpdjeJ4OpnBWhD/X8zT/iT1wYLiVU= -github.com/DataDog/datadog-go v0.0.0-20180330214955-e67964b4021a/go.mod h1:LButxg5PwREeZtORoXG3tL4fMGNddJ+vMq1mwgfaqoQ= github.com/DataDog/zstd v1.3.6-0.20190409195224-796139022798/go.mod h1:1jcaCB/ufaK+sKp1NBhlGmpz41jOoPQ35bpF36t7BBo= github.com/DataDog/zstd v1.4.0 h1:vhoV+DUHnRZdKW1i5UMjAk2G4JY8wN4ayRfYDNdEhwo= github.com/DataDog/zstd v1.4.0/go.mod h1:1jcaCB/ufaK+sKp1NBhlGmpz41jOoPQ35bpF36t7BBo= diff --git a/worker/create_batches_worker.go b/worker/create_batches_worker.go index d1c6c03..ac154a0 100644 --- a/worker/create_batches_worker.go +++ b/worker/create_batches_worker.go @@ -106,7 +106,7 @@ func (b *CreateBatchesWorker) getUserBatchFromPG(userIds *[]string, job *model.J start := time.Now() query := fmt.Sprintf("SELECT user_id, token, locale, tz FROM %s WHERE user_id IN (?)", GetPushDBTableName(job.App.Name, job.Service)) _, err := b.Workers.PushDB.Query(&users, query, pg.In(*userIds)) - b.Workers.Statsd.Timing("get_csv_batch_from_pg", time.Now().Sub(start), job.Labels(), 1) + observeWorkerDuration("get_csv_batch_from_pg", time.Since(start), job.Labels()) b.checkErr(job, err) return &users @@ -256,14 +256,14 @@ func (b *CreateBatchesWorker) Process(message *goworkers2.Msg) error { zap.Int("totalParts", msg.TotalParts), ) - b.Workers.Statsd.Incr(CreateBatchesWorkerStart, msg.Job.Labels(), 1) + incrWorkerEvent(CreateBatchesWorkerStart, msg.Job.Labels()) err = b.Workers.MarathonDB.Model(&msg.Job).Column("job.status").Relation("App").Where("job.id = ?", msg.Job.ID).Select() b.checkErr(&msg.Job, err) if msg.Job.Status == stoppedJobStatus { l.Info("stopped job") - b.Workers.Statsd.Incr(CreateBatchesWorkerCompleted, msg.Job.Labels(), 1) + incrWorkerEvent(CreateBatchesWorkerCompleted, msg.Job.Labels()) return nil } l.Info("starting") @@ -277,7 +277,7 @@ func (b *CreateBatchesWorker) Process(message *goworkers2.Msg) error { _, buffer, err := b.Workers.S3Client.DownloadChunk(int64(msg.Start), int64(msg.Size), msg.Job.CSVPath) labels := msg.Job.Labels() labels = append(labels, fmt.Sprintf("error:%t", err != nil)) - b.Workers.Statsd.Timing(GetCsvFromS3Timing, time.Now().Sub(start), labels, 1) + observeWorkerDuration(GetCsvFromS3Timing, time.Since(start), labels) b.checkErr(&msg.Job, err) ids := b.getIDs(buffer, &msg) @@ -297,10 +297,10 @@ func (b *CreateBatchesWorker) Process(message *goworkers2.Msg) error { b.checkErr(&msg.Job, err) //b.updateCompletedAt(time.Now().UnixNano(), &msg.Job) msg.Job.TagError(b.Workers.MarathonDB, nameCreateBatches, "the job has finished without finding any valid user ids") - b.Workers.Statsd.Incr(CreateBatchesWorkerError, msg.Job.Labels(), 1) + incrWorkerEvent(CreateBatchesWorkerError, msg.Job.Labels()) } else { msg.Job.TagSuccess(b.Workers.MarathonDB, nameCreateBatches, "finished") - b.Workers.Statsd.Incr(CreateBatchesWorkerCompleted, msg.Job.Labels(), 1) + incrWorkerEvent(CreateBatchesWorkerCompleted, msg.Job.Labels()) } // TODO: schedule a job to run after send all messages. This job will check @@ -319,7 +319,7 @@ func (b *CreateBatchesWorker) Process(message *goworkers2.Msg) error { func (b *CreateBatchesWorker) checkErr(job *model.Job, err error) { if err != nil { job.TagError(b.Workers.MarathonDB, nameCreateBatches, err.Error()) - b.Workers.Statsd.Incr(CreateBatchesWorkerError, job.Labels(), 1) + incrWorkerEvent(CreateBatchesWorkerError, job.Labels()) checkErr(b.Logger, err) } diff --git a/worker/csv_split.go b/worker/csv_split.go index 4e83762..fad771c 100644 --- a/worker/csv_split.go +++ b/worker/csv_split.go @@ -87,11 +87,11 @@ func (b *CSVSplitWorker) Process(message *goworkers2.Msg) error { l.Debug("job found") job.TagRunning(b.Workers.MarathonDB, nameSCVSplit, "starting") - b.Workers.Statsd.Incr(CsvSplitWorkerStart, job.Labels(), 1) + incrWorkerEvent(CsvSplitWorkerStart, job.Labels()) if job.Status == stoppedJobStatus { l.Info("stopped job") - b.Workers.Statsd.Incr(CsvSplitWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(CsvSplitWorkerCompleted, job.Labels()) return nil } @@ -117,11 +117,11 @@ func (b *CSVSplitWorker) Process(message *goworkers2.Msg) error { }) b.checkErr(job, err) start += size - b.Workers.Statsd.Incr("csv_job_part", job.Labels(), 1) + incrWorkerEvent("csv_job_part", job.Labels()) } job.TagSuccess(b.Workers.MarathonDB, nameSCVSplit, "finished") - b.Workers.Statsd.Incr(CsvSplitWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(CsvSplitWorkerCompleted, job.Labels()) l.Info("finished") return nil @@ -130,7 +130,7 @@ func (b *CSVSplitWorker) Process(message *goworkers2.Msg) error { func (b *CSVSplitWorker) checkErr(job *model.Job, err error) { if err != nil { job.TagError(b.Workers.MarathonDB, nameSCVSplit, err.Error()) - b.Workers.Statsd.Incr(CsvSplitWorkerError, job.Labels(), 1) + incrWorkerEvent(CsvSplitWorkerError, job.Labels()) checkErr(b.Logger, err) } diff --git a/worker/direct_worker.go b/worker/direct_worker.go index 95c50b3..9970fa6 100644 --- a/worker/direct_worker.go +++ b/worker/direct_worker.go @@ -127,26 +127,26 @@ func (b *DirectWorker) Process(message *goworkers2.Msg) error { job, err := b.Workers.GetJob(msg.JobUUID) checkErr(l, err) - b.Workers.Statsd.Incr(DirectWorkerStart, job.Labels(), 1) + incrWorkerEvent(DirectWorkerStart, job.Labels()) if job.ExpiresAt > 0 && job.ExpiresAt < time.Now().UnixNano() { log.I(l, "expired") - b.Workers.Statsd.Incr(DirectWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(DirectWorkerCompleted, job.Labels()) return nil } switch job.Status { case "circuitbreak": log.I(l, "circuit break") - b.Workers.Statsd.Incr(DirectWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(DirectWorkerCompleted, job.Labels()) return nil case "paused": log.I(l, "paused") - b.Workers.Statsd.Incr(DirectWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(DirectWorkerCompleted, job.Labels()) return nil case "stopped": log.I(l, "stopped") - b.Workers.Statsd.Incr(DirectWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(DirectWorkerCompleted, job.Labels()) return nil default: log.D(l, "valid") @@ -168,7 +168,7 @@ func (b *DirectWorker) Process(message *goworkers2.Msg) error { l.Error("Error fetching users", zap.Error(err)) } - b.Workers.Statsd.Timing(GetUsersFromDbTiming, time.Now().Sub(start), job.Labels(), 1) + observeWorkerDuration(GetUsersFromDbTiming, time.Since(start), job.Labels()) successfulUsers := len(users) @@ -271,7 +271,7 @@ func (b *DirectWorker) Process(message *goworkers2.Msg) error { _, err = b.Workers.ScheduleJobCompletedJob(job.ID.String(), at) } - b.Workers.Statsd.Incr(DirectWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(DirectWorkerCompleted, job.Labels()) l.Info("finished") return nil @@ -280,7 +280,7 @@ func (b *DirectWorker) Process(message *goworkers2.Msg) error { func (b *DirectWorker) checkErr(job *model.Job, err error) { if err != nil { job.TagError(b.Workers.MarathonDB, nameDirectWorker, err.Error()) - b.Workers.Statsd.Incr(DirectWorkerError, job.Labels(), 1) + incrWorkerEvent(DirectWorkerError, job.Labels()) checkErr(b.Logger, err) } diff --git a/worker/job_completed_worker.go b/worker/job_completed_worker.go index 9ccf224..106441f 100644 --- a/worker/job_completed_worker.go +++ b/worker/job_completed_worker.go @@ -101,7 +101,7 @@ func (b *JobCompletedWorker) Process(message *goworkers2.Msg) error { job, err := b.Workers.GetJob(id) checkErr(l, err) - b.Workers.Statsd.Incr(JobCompletedWorkerStart, job.Labels(), 1) + incrWorkerEvent(JobCompletedWorkerStart, job.Labels()) job.TagRunning(b.Workers.MarathonDB, nameJobCompleted, "starting") @@ -114,7 +114,7 @@ func (b *JobCompletedWorker) Process(message *goworkers2.Msg) error { b.flushControlGroup(job) job.TagSuccess(b.Workers.MarathonDB, nameJobCompleted, "finished") - b.Workers.Statsd.Incr(JobCompletedWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(JobCompletedWorkerCompleted, job.Labels()) log.I(l, "finished") @@ -124,7 +124,7 @@ func (b *JobCompletedWorker) Process(message *goworkers2.Msg) error { func (b *JobCompletedWorker) checkErr(job *model.Job, err error) { if err != nil { job.TagError(b.Workers.MarathonDB, nameJobCompleted, err.Error()) - b.Workers.Statsd.Incr(JobCompletedWorkerError, job.Labels(), 1) + incrWorkerEvent(JobCompletedWorkerError, job.Labels()) checkErr(b.Logger, err) } diff --git a/worker/metrics.go b/worker/metrics.go index f0e9c50..cb180a4 100644 --- a/worker/metrics.go +++ b/worker/metrics.go @@ -1,5 +1,12 @@ package worker +import ( + "sync" + "time" + + "github.com/prometheus/client_golang/prometheus" +) + const ( CreateBatchesWorkerStart = "starting_create_batches_worker" CreateBatchesWorkerCompleted = "completed_create_batches_worker" @@ -28,3 +35,62 @@ const ( GetCsvFromS3Timing = "get_csv_from_s3" GetUsersFromDbTiming = "get_from_pg" ) + +var ( + registerMetricsOnce sync.Once + + // workerCounter counts worker lifecycle events (start/completed/error and csv_job_part). + // Label names match the statsd tag keys previously passed via job.Labels(): + // game=, platform= + // The metric name uses underscores; the DD OpenMetrics check maps marathon_* → marathon.* + workerCounter = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "marathon_worker_events_total", + Help: "Worker lifecycle event counter (start/completed/error) per worker type and job labels.", + }, + []string{"event", "game", "platform"}, + ) + + // workerDuration observes timing operations (get_csv_from_s3, get_from_pg, get_csv_batch_from_pg, save_control_group). + workerDuration = prometheus.NewHistogramVec( + prometheus.HistogramOpts{ + Name: "marathon_worker_duration_milliseconds", + Help: "Worker operation duration in milliseconds.", + Buckets: []float64{5, 10, 25, 50, 100, 250, 500, 1000, 2500, 5000, 10000, 30000}, + }, + []string{"operation", "game", "platform"}, + ) +) + +// MustRegisterMetrics registers Marathon worker Prometheus collectors. Safe to call more than once. +func MustRegisterMetrics() { + registerMetricsOnce.Do(func() { + prometheus.MustRegister(workerCounter) + prometheus.MustRegister(workerDuration) + }) +} + +// incrWorkerEvent increments the worker event counter for the given event name and job labels. +// labels is the slice returned by job.Labels() — each element has the form "key:value". +func incrWorkerEvent(event string, labels []string) { + game, platform := parseJobLabels(labels) + workerCounter.WithLabelValues(event, game, platform).Inc() +} + +// observeWorkerDuration records an operation duration for the given timing name and job labels. +func observeWorkerDuration(operation string, d time.Duration, labels []string) { + game, platform := parseJobLabels(labels) + workerDuration.WithLabelValues(operation, game, platform).Observe(float64(d.Milliseconds())) +} + +// parseJobLabels extracts game and platform from a Labels() slice of "key:value" strings. +func parseJobLabels(labels []string) (game, platform string) { + for _, l := range labels { + if len(l) > 5 && l[:5] == "game:" { + game = l[5:] + } else if len(l) > 9 && l[:9] == "platform:" { + platform = l[9:] + } + } + return +} diff --git a/worker/metrics_test.go b/worker/metrics_test.go new file mode 100644 index 0000000..27f2394 --- /dev/null +++ b/worker/metrics_test.go @@ -0,0 +1,99 @@ +package worker_test + +import ( + "testing" + + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" + "github.com/topfreegames/marathon/extensions" +) + +// collectCounterValue gathers metrics from a registry and returns the value of the +// counter series identified by the given metric family name and label pairs. +func collectCounterValue(t *testing.T, reg *prometheus.Registry, metricName string, labelPairs map[string]string) float64 { + t.Helper() + mfs, err := reg.Gather() + if err != nil { + t.Fatalf("gather: %v", err) + } + for _, mf := range mfs { + if mf.GetName() != metricName { + continue + } + for _, m := range mf.GetMetric() { + if labelsMatch(m.GetLabel(), labelPairs) { + return m.GetCounter().GetValue() + } + } + } + return 0 +} + +func labelsMatch(got []*dto.LabelPair, want map[string]string) bool { + matched := 0 + for _, lp := range got { + if v, ok := want[lp.GetName()]; ok && v == lp.GetValue() { + matched++ + } + } + return matched == len(want) +} + +func TestWorkerEventCounter(t *testing.T) { + reg := prometheus.NewRegistry() + + counter := prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "marathon_worker_events_total", + Help: "Worker lifecycle event counter.", + }, + []string{"event", "game", "platform"}, + ) + reg.MustRegister(counter) + + // Simulate incrWorkerEvent for a worker start event. + counter.WithLabelValues("starting_create_batches_worker", "mygame", "gcm").Inc() + + val := collectCounterValue(t, reg, "marathon_worker_events_total", map[string]string{ + "event": "starting_create_batches_worker", + "game": "mygame", + "platform": "gcm", + }) + if val <= 0 { + t.Fatalf("expected counter > 0, got %v", val) + } +} + +func TestKafkaSendMessageReturnCounter(t *testing.T) { + reg := prometheus.NewRegistry() + + counter := prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "marathon_kafka_send_message_return_total", + Help: "Kafka async producer message acknowledgements.", + }, + []string{"error"}, + ) + reg.MustRegister(counter) + + // Simulate the success goroutine in connectToKafka. + counter.WithLabelValues("false").Inc() + counter.WithLabelValues("false").Inc() + counter.WithLabelValues("true").Inc() + + successVal := collectCounterValue(t, reg, "marathon_kafka_send_message_return_total", map[string]string{"error": "false"}) + if successVal != 2 { + t.Fatalf("expected success counter = 2, got %v", successVal) + } + errorVal := collectCounterValue(t, reg, "marathon_kafka_send_message_return_total", map[string]string{"error": "true"}) + if errorVal != 1 { + t.Fatalf("expected error counter = 1, got %v", errorVal) + } +} + +// TestKafkaMetricsRegistration verifies MustRegisterKafkaMetrics is safe to call repeatedly. +func TestKafkaMetricsRegistration(t *testing.T) { + // Must not panic on multiple calls (sync.Once guard). + extensions.MustRegisterKafkaMetrics() + extensions.MustRegisterKafkaMetrics() +} diff --git a/worker/process_batch_worker.go b/worker/process_batch_worker.go index 9fadc07..f8cf358 100644 --- a/worker/process_batch_worker.go +++ b/worker/process_batch_worker.go @@ -185,11 +185,11 @@ func (b *ProcessBatchWorker) Process(message *goworkers2.Msg) error { ) log.D(l, "Retrieved job successfully.") - b.Workers.Statsd.Incr(ProcessBatchWorkerStart, job.Labels(), 1) + incrWorkerEvent(ProcessBatchWorkerStart, job.Labels()) if job.ExpiresAt > 0 && job.ExpiresAt < time.Now().UnixNano() { log.I(l, "expired") - b.Workers.Statsd.Incr(ProcessBatchWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(ProcessBatchWorkerCompleted, job.Labels()) return nil } @@ -197,16 +197,16 @@ func (b *ProcessBatchWorker) Process(message *goworkers2.Msg) error { case "circuitbreak": log.I(l, "circuit break") b.moveJobToPausedQueue(job, message) - b.Workers.Statsd.Incr(ProcessBatchWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(ProcessBatchWorkerCompleted, job.Labels()) return nil case "paused": log.I(l, "paused") b.moveJobToPausedQueue(job, message) - b.Workers.Statsd.Incr(ProcessBatchWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(ProcessBatchWorkerCompleted, job.Labels()) return nil case "stopped": log.I(l, "stopped") - b.Workers.Statsd.Incr(ProcessBatchWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(ProcessBatchWorkerCompleted, job.Labels()) return nil default: log.D(l, "valid") @@ -303,7 +303,7 @@ func (b *ProcessBatchWorker) Process(message *goworkers2.Msg) error { b.checkErr(job, fmt.Errorf("failed to send message to several users, considering batch as failed")) } - b.Workers.Statsd.Incr(ProcessBatchWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(ProcessBatchWorkerCompleted, job.Labels()) log.I(l, "finished") return nil @@ -312,7 +312,7 @@ func (b *ProcessBatchWorker) Process(message *goworkers2.Msg) error { func (b *ProcessBatchWorker) checkErr(job *model.Job, err error) { if err != nil { job.TagError(b.Workers.MarathonDB, nameProcessBatchWorker, err.Error()) - b.Workers.Statsd.Incr(ProcessBatchWorkerError, job.Labels(), 1) + incrWorkerEvent(ProcessBatchWorkerError, job.Labels()) checkErr(b.Logger, err) } diff --git a/worker/resume_job_worker.go b/worker/resume_job_worker.go index 11463ff..486f8e5 100644 --- a/worker/resume_job_worker.go +++ b/worker/resume_job_worker.go @@ -66,14 +66,14 @@ func (b *ResumeJobWorker) Process(message *goworkers2.Msg) error { job, err := b.Workers.GetJob(id) checkErr(l, err) - b.Workers.Statsd.Incr(ResumeJobWorkerStart, job.Labels(), 1) + incrWorkerEvent(ResumeJobWorkerStart, job.Labels()) if job.Status == stoppedJobStatus { l.Info("stopped job resume_job_worker") err := b.Workers.RedisClient.Del(fmt.Sprintf("%s-pausedjobs", jobID.(string))).Err() if err != nil && err != redis.Nil { checkErr(b.Logger, err) } - b.Workers.Statsd.Incr(ResumeJobWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(ResumeJobWorkerCompleted, job.Labels()) return nil } @@ -93,7 +93,7 @@ func (b *ResumeJobWorker) Process(message *goworkers2.Msg) error { b.checkErr(job, err) } - b.Workers.Statsd.Incr(ResumeJobWorkerCompleted, job.Labels(), 1) + incrWorkerEvent(ResumeJobWorkerCompleted, job.Labels()) log.I(b.Logger, "finished resume_job_worker") return nil @@ -102,7 +102,7 @@ func (b *ResumeJobWorker) Process(message *goworkers2.Msg) error { func (b *ResumeJobWorker) checkErr(job *model.Job, err error) { if err != nil { job.TagError(b.Workers.MarathonDB, ResumeJobWorkerError, err.Error()) - b.Workers.Statsd.Incr(ResumeJobWorkerError, job.Labels(), 1) + incrWorkerEvent(ResumeJobWorkerError, job.Labels()) checkErr(b.Logger, err) } diff --git a/worker/worker.go b/worker/worker.go index a9d19ba..c264e0d 100644 --- a/worker/worker.go +++ b/worker/worker.go @@ -31,10 +31,10 @@ import ( "strings" "time" - "github.com/DataDog/datadog-go/statsd" goworkers2 "github.com/digitalocean/go-workers2" raven "github.com/getsentry/raven-go" pg "github.com/go-pg/pg/v10" + "github.com/prometheus/client_golang/prometheus/promhttp" uuid "github.com/satori/go.uuid" "github.com/spf13/viper" "github.com/topfreegames/marathon/extensions" @@ -53,7 +53,6 @@ type Worker struct { DBPageSize int S3Client interfaces.S3 PageProcessingConcurrency int - Statsd *statsd.Client RedisClient *redis.Client ConfigPath string SendgridClient *extensions.SendgridClient @@ -89,14 +88,14 @@ func (w *Worker) configure() { w.loadConfigurationDefaults() w.configureSentry() w.configureRedis() - w.configureStatsd() w.configureWorkers() - w.configureStatsd() w.configurePushDatabase() w.configureMarathonDatabase() w.configureS3Client() w.configureSendgrid() w.configureKafkaProducer() + MustRegisterMetrics() + extensions.MustRegisterKafkaMetrics() } func (w *Worker) loadConfigurationDefaults() { @@ -104,10 +103,9 @@ func (w *Worker) loadConfigurationDefaults() { w.Config.SetDefault("workers.redis.database", "0") w.Config.SetDefault("workers.redis.poolSize", "10") w.Config.SetDefault("workers.statsPort", 8081) + w.Config.SetDefault("workers.metricsPort", 9090) w.Config.SetDefault("workers.concurrency", 10) w.Config.SetDefault("database.url", "postgres://localhost:5432/marathon?sslmode=disable") - w.Config.SetDefault("workers.statsd.host", "127.0.0.1:8125") - w.Config.SetDefault("workers.statsd.prefix", "marathon.") } func (w *Worker) configureSendgrid() { @@ -129,18 +127,6 @@ func (w *Worker) configureMarathonDatabase() { w.MarathonDB = connection.DB } -func (w *Worker) configureStatsd() { - host := w.Config.GetString("workers.statsd.host") - prefix := w.Config.GetString("workers.statsd.prefix") - - client, err := statsd.New(host) - if err != nil { - return - } - client.Namespace = prefix - w.Statsd = client -} - func (w *Worker) configureRedis() { redisHost := w.Config.GetString("workers.redis.host") redisPort := w.Config.GetInt("workers.redis.port") @@ -230,7 +216,7 @@ func (w *Worker) configureSentry() { func (w *Worker) configureKafkaProducer() { var kafka *extensions.KafkaProducer var err error - kafka, err = extensions.NewKafkaProducer(w.Config, w.Logger, w.Statsd) + kafka, err = extensions.NewKafkaProducer(w.Config, w.Logger) checkErr(w.Logger, err) w.Kafka = kafka } @@ -434,9 +420,10 @@ func (w *Worker) ScheduleJobCompletedJob(jobID string, at int64) (string, error) // Start starts the worker func (w *Worker) Start() { jobsStatsPort := w.Config.GetInt("workers.statsPort") + metricsPort := w.Config.GetInt("workers.metricsPort") go func() { - http.HandleFunc("/stats", func(rw http.ResponseWriter, req *http.Request) { - + mux := http.NewServeMux() + mux.HandleFunc("/stats", func(rw http.ResponseWriter, req *http.Request) { _, marathonError := w.MarathonDB.Exec("SELECT 1") _, pushError := w.PushDB.Exec("SELECT 1") pong, redisError := w.RedisClient.Ping().Result() @@ -456,7 +443,14 @@ func (w *Worker) Start() { } json.NewEncoder(rw).Encode(status) }) - if err := http.ListenAndServe(fmt.Sprint(":", jobsStatsPort), nil); err != nil { + if err := http.ListenAndServe(fmt.Sprint(":", jobsStatsPort), mux); err != nil { + panic(err) + } + }() + go func() { + metricsMux := http.NewServeMux() + metricsMux.Handle("/metrics", promhttp.Handler()) + if err := http.ListenAndServe(fmt.Sprint(":", metricsPort), metricsMux); err != nil { panic(err) } }() @@ -472,7 +466,7 @@ func (w *Worker) SendControlGroupToRedis(job *model.Job, ids []string) { args = append(args, id) } w.RedisClient.LPush(fmt.Sprintf("%s-CONTROL", hash), args...).Result() - w.Statsd.Timing("save_control_group", time.Now().Sub(start), job.Labels(), 1) + observeWorkerDuration("save_control_group", time.Since(start), job.Labels()) } // GetJob get a job from the db