Skip to content

Commit 536c51e

Browse files
committed
reduce unneccessary usage of sync.atomic
pprof shows it might have a larger impact than thought
1 parent 760ff8e commit 536c51e

2 files changed

Lines changed: 47 additions & 27 deletions

File tree

message_pool.go

Lines changed: 43 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,6 @@ import (
99
"runtime/debug"
1010
"strings"
1111
"sync"
12-
"sync/atomic"
1312
"time"
1413

1514
"golang.org/x/exp/maps"
@@ -42,16 +41,15 @@ const MessagePoolFlagShared = uint8(0x01)
4241

4342
var InitialMessagePoolByteCount = mib(1)
4443

45-
var nextId atomic.Uint64
46-
4744
type messagePool struct {
4845
size int
4946
pool [][]byte
5047
count int
51-
takenTags [256]atomic.Uint64
52-
returnedTags [256]atomic.Uint64
53-
createdTags [256]atomic.Uint64
48+
takenTags [256]uint64
49+
returnedTags [256]uint64
50+
createdTags [256]uint64
5451
stateLock sync.Mutex
52+
nextId uint64
5553
}
5654

5755
func newMessagePool(size int, maxCount int) *messagePool {
@@ -98,7 +96,8 @@ func (self *messagePool) Get() []byte {
9896

9997
// create a new message
10098
poolMessage := make([]byte, self.size+MessagePoolMetaByteCount)
101-
binary.BigEndian.PutUint64(poolMessage[self.size:], nextId.Add(1))
99+
self.nextId += 1
100+
binary.BigEndian.PutUint64(poolMessage[self.size:], self.nextId)
102101
poolMessage[self.size+8] = 255
103102
return poolMessage
104103
}
@@ -133,10 +132,17 @@ func poolStats(pools []*messagePool) {
133132
for {
134133
for _, pool := range pools {
135134
for tag := range 256 {
136-
taken := pool.takenTags[tag].Load()
135+
var taken uint64
136+
var returned uint64
137+
var created uint64
138+
func() {
139+
pool.stateLock.Lock()
140+
defer pool.stateLock.Unlock()
141+
taken = pool.takenTags[tag]
142+
returned = pool.returnedTags[tag]
143+
created = pool.createdTags[tag]
144+
}()
137145
if 0 < taken {
138-
returned := pool.returnedTags[tag].Load()
139-
created := pool.createdTags[tag].Load()
140146
ratio := float32(returned) / float32(taken)
141147
reuse := float32(taken-created) / float32(taken)
142148
var caller string
@@ -200,11 +206,15 @@ func debugTag() uint8 {
200206

201207
func ResetMessagePoolStats() {
202208
for _, pool := range orderedMessagePools() {
203-
for tag := range 256 {
204-
pool.takenTags[tag].Store(0)
205-
pool.returnedTags[tag].Store(0)
206-
pool.createdTags[tag].Store(0)
207-
}
209+
func() {
210+
pool.stateLock.Lock()
211+
defer pool.stateLock.Unlock()
212+
for tag := range 256 {
213+
pool.takenTags[tag] = 0
214+
pool.returnedTags[tag] = 0
215+
pool.createdTags[tag] = 0
216+
}
217+
}()
208218
}
209219
}
210220

@@ -213,9 +223,16 @@ func MessagePoolStats() map[int]map[int]float32 {
213223
for _, pool := range orderedMessagePools() {
214224
tagRatios := map[int]float32{}
215225
for tag := range 256 {
216-
taken := pool.takenTags[tag].Load()
226+
var taken uint64
227+
var returned uint64
228+
func() {
229+
pool.stateLock.Lock()
230+
defer pool.stateLock.Unlock()
231+
taken = pool.takenTags[tag]
232+
returned = pool.returnedTags[tag]
233+
}()
234+
217235
if 0 < taken {
218-
returned := pool.returnedTags[tag].Load()
219236
ratio := float32(returned) / float32(taken)
220237
tagRatios[tag] = ratio
221238
}
@@ -323,17 +340,19 @@ func MessagePoolGetDetailedWithTag(n int, tag uint8) ([]byte, bool) {
323340
for _, pool := range orderedMessagePools {
324341
if n <= pool.size {
325342
poolMessage := pool.Get()
326-
if poolMessage[pool.size+8] == 255 {
327-
pool.createdTags[tag].Add(1)
328-
}
343+
c := poolMessage[pool.size+8] == 255
329344
poolMessage[pool.size+8] = tag
330-
pool.takenTags[tag].Add(1)
331345
id := binary.BigEndian.Uint64(poolMessage[pool.size:])
332346

333347
func() {
334348
pool.stateLock.Lock()
335349
defer pool.stateLock.Unlock()
336350

351+
if c {
352+
pool.createdTags[tag] += 1
353+
}
354+
pool.takenTags[tag] += 1
355+
337356
count := binary.BigEndian.Uint16(poolMessage[pool.size+10:])
338357

339358
if count != 0 {
@@ -362,6 +381,8 @@ func MessagePoolReturn(message []byte) bool {
362381
poolMessage := message[:c]
363382
id := binary.BigEndian.Uint64(poolMessage[pool.size:])
364383

384+
tag := poolMessage[pool.size+8]
385+
365386
r := false
366387
func() {
367388
pool.stateLock.Lock()
@@ -375,18 +396,17 @@ func MessagePoolReturn(message []byte) bool {
375396
}
376397
} else if count == 1 {
377398
r = true
399+
pool.returnedTags[tag] += 1
378400
} else {
379401
binary.BigEndian.PutUint16(poolMessage[pool.size+10:], count-1)
380402
}
381403
}()
382404

383405
if r {
384-
tag := poolMessage[pool.size+8]
385406
poolMessage[pool.size+8] = 0
386407
poolMessage[pool.size+9] = 0
387408
binary.BigEndian.PutUint16(poolMessage[pool.size+10:], 0)
388409
pool.Put(poolMessage)
389-
pool.returnedTags[tag].Add(1)
390410
return true
391411
}
392412
return false

transport.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -498,8 +498,8 @@ func (self *PlatformTransport) runH1(initialTimeout time.Duration) {
498498
handleCtx, handleCancel := context.WithCancel(self.ctx)
499499
defer handleCancel()
500500

501-
readCounter := atomic.Uint64{}
502-
writeCounter := atomic.Uint64{}
501+
var readCounter atomic.Uint64
502+
var writeCounter atomic.Uint64
503503

504504
send := make(chan []byte, self.settings.TransportBufferSize)
505505
receive := make(chan []byte, self.settings.TransportBufferSize)
@@ -1040,8 +1040,8 @@ func (self *PlatformTransport) runH3(ptMode TransportMode, initialTimeout time.D
10401040

10411041
framer := NewFramer(self.settings.FramerSettings)
10421042

1043-
readCounter := atomic.Uint64{}
1044-
writeCounter := atomic.Uint64{}
1043+
var readCounter atomic.Uint64
1044+
var writeCounter atomic.Uint64
10451045

10461046
send := make(chan []byte, self.settings.TransportBufferSize)
10471047
receive := make(chan []byte, self.settings.TransportBufferSize)

0 commit comments

Comments
 (0)