/
main.go
84 lines (76 loc) · 2.52 KB
/
main.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
package consumer
import (
"fmt"
"log"
"math"
"math/rand"
"net/http"
"strconv"
"sync"
"time"
"github.com/go-redis/redis"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/lwolf/konsumerator/hack/faker/lib"
)
var (
consumptionOffsetMetric = prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: "konsumerator",
Name: "messages_consumption_offset",
Help: "Last seen offset per partition",
}, []string{"partition"})
)
func runServer(port int) {
log.Printf("starting web server on port %d", port)
http.Handle("/metrics", promhttp.Handler())
log.Fatal(http.ListenAndServe(fmt.Sprintf(":%d", port), nil))
}
func consume(partition int, ratePerCore int) float64 {
milliCores, err := lib.GetCpuRequests()
if err != nil {
panic(fmt.Sprintf("unable to get cgroups value %v", err))
}
rand.Seed(int64(partition))
consumptionRate := (float64(milliCores) / 1000) * float64(ratePerCore)
fuzz := rand.Float64() * consumptionRate * 0.1
return fuzz + consumptionRate
}
func runConsumer(client *redis.Client, partition int, ratePerCore int) {
log.Printf("starting consumer for partition %d", partition)
productionOffsetKey := fmt.Sprintf("%s-%d", lib.ProductionOffsetKey, partition)
consumptionOffsetKey := fmt.Sprintf("%s-%d", lib.ConsumptionOffsetKey, partition)
ticker := time.NewTicker(time.Second)
for range ticker.C {
cOffset, err := lib.GetOffset(client, consumptionOffsetKey, 0)
if err != nil {
log.Fatalf("unable to load state from redis %v", err)
}
pOffset, err := lib.GetOffset(client, productionOffsetKey, 0)
if err != nil {
log.Fatalf("unable to get prod offset from redis %v", err)
}
proposedBatch := consume(partition, ratePerCore)
actualBatch := math.Max(0, math.Min(proposedBatch, float64(pOffset-cOffset)))
state := cOffset + int(actualBatch)
err = lib.SetOffset(client, consumptionOffsetKey, state)
if err != nil {
log.Printf("unable to set offset %v", err)
continue
}
consumptionOffsetMetric.WithLabelValues(strconv.Itoa(partition)).Set(float64(state))
log.Printf("processed %d(%d) messages from partition %d, new offset=%d", int(actualBatch), int(proposedBatch), partition, state)
}
}
func RunConsumer(redisClient *redis.Client, partitions []int, ratePerCore int, port int) {
prometheus.MustRegister(consumptionOffsetMetric)
wg := sync.WaitGroup{}
for _, p := range partitions {
wg.Add(1)
go func(p int) {
runConsumer(redisClient, p, ratePerCore)
wg.Done()
}(p)
}
runServer(port)
wg.Wait()
}