i've caused myself pretty weird issues with implementing atomic counter in my main event loop that is processing kafka msgs.. (needed quick&dirty solution to set some kind of a limiter in the event handler (kafka) for the items being processed in parallel, so when it hits 50 it will block getting a new event from kafka)
problems I'm observing are for example msg are not commited so they are delivered by kafka multiple times, also from time to time event processing ends up with the msg 'kafka event read' and it just doesn't enter processMessage()
any idea what am I doing wrong? (fyi EventStream is a chan of event structs)
MaxConcurReqs = 50
var reqCounter int32 // this is how I want to control amount of requests that can be processed concurrently
for event := range eh.EventStream() {
if event.StreamError != nil {
log.WithField("error", event.StreamError.Error()).Warning("msg stream error")
continue
}
log.WithField("event", string(event.Message.Payload())).Debug("kafka event read")
CountCheck:
for {
if atomic.LoadInt32(&reqCounter) < int32(config.MaxConcurReqs) {
atomic.AddInt32(&reqCounter, 1)
break CountCheck
}
log.Debugf("maximum concurrent requests are in progress, waiting for running jobs to finish",
atomic.LoadInt32(&reqCounter),
config.MaxConcurReqs,
)
time.Sleep(1 * time.Second)
}
go func() {
processMessage(event.Message.Context(), event.Message.Payload(), config, handlers)
atomic.AddInt32(&reqCounter, -1) // remove finished job from the 'queue'
}()
if err := event.Message.Commit(); err != nil {
log.WithFields(log.Fields{
"event": string(event.Message.Payload()),
"error": err.Error(),
}).Error("kafka commit error")
}
}