Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 24 additions & 6 deletions extensions/kafka_producer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()
Expand Down Expand Up @@ -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)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new metric is called marathon_kafka_send_message_return_total, is the renaming expected? If that's the case, let's add a document under docs/ listing the changes in metrics and what consumers must do to adapt

kafkaSendMessageReturn.WithLabelValues("false").Inc()
}
}()

go func() {
for range producer.Errors() {
c.Statsd.Incr("send_message_return", []string{"error:true"}, 1)
kafkaSendMessageReturn.WithLabelValues("true").Inc()
}
}()

Expand Down
11 changes: 3 additions & 8 deletions extensions/kafka_producer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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(
Expand All @@ -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() {
Expand All @@ -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()

Expand All @@ -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()

Expand All @@ -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()

Expand Down
5 changes: 2 additions & 3 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
2 changes: 0 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
14 changes: 7 additions & 7 deletions worker/create_batches_worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand All @@ -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)
Expand All @@ -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
Expand All @@ -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)
}
Expand Down
10 changes: 5 additions & 5 deletions worker/csv_split.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand All @@ -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
Expand All @@ -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)
}
Expand Down
16 changes: 8 additions & 8 deletions worker/direct_worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -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)

Expand Down Expand Up @@ -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
Expand All @@ -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)
}
Expand Down
6 changes: 3 additions & 3 deletions worker/job_completed_worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand All @@ -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")

Expand All @@ -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)
}
Expand Down
Loading
Loading