#Atomic counter (limiter)

3 messages · Page 1 of 1 (latest)

unique parcel
#

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")
        }
    }
leaden dirge
#

Hey, Kafka is pretty tricky to get right. Dunno if I can help you with that code, but, if you're open to using a higher level framework, I can suggest Benthos. It has a builtin configurable rate limiter. You can use it as a CLI, but you can also import it as an API and use its kafka_franz input and write your own custom output or even read the messages it fetches in some custom code by tapping directly into its internal event stream: https://pkg.go.dev/github.com/benthosdev/benthos/[email protected]/public/service#StreamBuilder.AddConsumerFunc

unique parcel
#

I'll take a look, thanks