Skip to content

Commit

Permalink
fix: data race in roundrobin balancer (#1251)
Browse files Browse the repository at this point in the history
  • Loading branch information
petedannemann authored Dec 13, 2023
1 parent edf45f6 commit 2af3101
Showing 1 changed file with 8 additions and 4 deletions.
12 changes: 8 additions & 4 deletions balancer.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import (
"math/rand"
"sort"
"sync"
"sync/atomic"
)

// The Balancer interface provides an abstraction of the message distribution
Expand Down Expand Up @@ -42,8 +41,10 @@ func (f BalancerFunc) Balance(msg Message, partitions ...int) int {
type RoundRobin struct {
ChunkSize int
// Use a 32 bits integer so RoundRobin values don't need to be aligned to
// apply atomic increments.
// apply increments.
counter uint32

mutex sync.Mutex
}

// Balance satisfies the Balancer interface.
Expand All @@ -52,14 +53,17 @@ func (rr *RoundRobin) Balance(msg Message, partitions ...int) int {
}

func (rr *RoundRobin) balance(partitions []int) int {
rr.mutex.Lock()
defer rr.mutex.Unlock()

if rr.ChunkSize < 1 {
rr.ChunkSize = 1
}

length := len(partitions)
counterNow := atomic.LoadUint32(&rr.counter)
counterNow := rr.counter
offset := int(counterNow / uint32(rr.ChunkSize))
atomic.AddUint32(&rr.counter, 1)
rr.counter++
return partitions[offset%length]
}

Expand Down

0 comments on commit 2af3101

Please sign in to comment.