为高度并发的应用程序实现全局计数器的最佳方法是什么?在我的情况下,我可能有10K-20K的例程执行"工作",我想计算例程正在集体工作的项目的数量和类型...
"经典"同步编码风格如下所示:
var work_counter int
func GoWorkerRoutine() {
for {
// do work
atomic.AddInt32(&work_counter,1)
}
}
Run Code Online (Sandbox Code Playgroud)
现在这变得更加复杂,因为我想跟踪正在完成的工作的"类型",所以我真的需要这样的东西:
var work_counter map[string]int
var work_mux sync.Mutex
func GoWorkerRoutine() {
for {
// do work
work_mux.Lock()
work_counter["type1"]++
work_mux.Unlock()
}
}
Run Code Online (Sandbox Code Playgroud)
似乎应该使用渠道或类似的东西"go"优化方式:
var work_counter int
var work_chan chan int // make() called somewhere else (buffered)
// started somewher else
func GoCounterRoutine() {
for {
select {
case c := <- work_chan:
work_counter += c
break
}
}
}
func GoWorkerRoutine() {
for {
// do work
work_chan <- 1
}
}
Run Code Online (Sandbox Code Playgroud)
最后一个例子仍然缺少地图,但这很容易添加.这种风格是否会提供比简单的原子增量更好的性能?当我们谈论对全局值的并发访问与可能阻止I/O完成的事情时,我无法判断这是否或多或少复杂......
我们很感激.
2013年5月28日更新:
我测试了几个实现,结果不是我的预期,这是我的计数器源代码:
package helpers
import (
)
type CounterIncrementStruct struct {
bucket string
value int
}
type CounterQueryStruct struct {
bucket string
channel chan int
}
var counter map[string]int
var counterIncrementChan chan CounterIncrementStruct
var counterQueryChan chan CounterQueryStruct
var counterListChan chan chan map[string]int
func CounterInitialize() {
counter = make(map[string]int)
counterIncrementChan = make(chan CounterIncrementStruct,0)
counterQueryChan = make(chan CounterQueryStruct,100)
counterListChan = make(chan chan map[string]int,100)
go goCounterWriter()
}
func goCounterWriter() {
for {
select {
case ci := <- counterIncrementChan:
if len(ci.bucket)==0 { return }
counter[ci.bucket]+=ci.value
break
case cq := <- counterQueryChan:
val,found:=counter[cq.bucket]
if found {
cq.channel <- val
} else {
cq.channel <- -1
}
break
case cl := <- counterListChan:
nm := make(map[string]int)
for k, v := range counter {
nm[k] = v
}
cl <- nm
break
}
}
}
func CounterIncrement(bucket string, counter int) {
if len(bucket)==0 || counter==0 { return }
counterIncrementChan <- CounterIncrementStruct{bucket,counter}
}
func CounterQuery(bucket string) int {
if len(bucket)==0 { return -1 }
reply := make(chan int)
counterQueryChan <- CounterQueryStruct{bucket,reply}
return <- reply
}
func CounterList() map[string]int {
reply := make(chan map[string]int)
counterListChan <- reply
return <- reply
}
Run Code Online (Sandbox Code Playgroud)
它使用通道进行写入和读取,这似乎是合乎逻辑的.
这是我的测试用例:
func bcRoutine(b *testing.B,e chan bool) {
for i := 0; i < b.N; i++ {
CounterIncrement("abc123",5)
CounterIncrement("def456",5)
CounterIncrement("ghi789",5)
CounterIncrement("abc123",5)
CounterIncrement("def456",5)
CounterIncrement("ghi789",5)
}
e<-true
}
func BenchmarkChannels(b *testing.B) {
b.StopTimer()
CounterInitialize()
e:=make(chan bool)
b.StartTimer()
go bcRoutine(b,e)
go bcRoutine(b,e)
go bcRoutine(b,e)
go bcRoutine(b,e)
go bcRoutine(b,e)
<-e
<-e
<-e
<-e
<-e
}
var mux sync.Mutex
var m map[string]int
func bmIncrement(bucket string, value int) {
mux.Lock()
m[bucket]+=value
mux.Unlock()
}
func bmRoutine(b *testing.B,e chan bool) {
for i := 0; i < b.N; i++ {
bmIncrement("abc123",5)
bmIncrement("def456",5)
bmIncrement("ghi789",5)
bmIncrement("abc123",5)
bmIncrement("def456",5)
bmIncrement("ghi789",5)
}
e<-true
}
func BenchmarkMutex(b *testing.B) {
b.StopTimer()
m=make(map[string]int)
e:=make(chan bool)
b.StartTimer()
for i := 0; i < b.N; i++ {
bmIncrement("abc123",5)
bmIncrement("def456",5)
bmIncrement("ghi789",5)
bmIncrement("abc123",5)
bmIncrement("def456",5)
bmIncrement("ghi789",5)
}
go bmRoutine(b,e)
go bmRoutine(b,e)
go bmRoutine(b,e)
go bmRoutine(b,e)
go bmRoutine(b,e)
<-e
<-e
<-e
<-e
<-e
}
Run Code Online (Sandbox Code Playgroud)
我实现了一个简单的基准测试,在地图周围只有一个互斥体(只是测试写入),并使用并行运行的5个goroutine进行基准测试.结果如下:
$ go test --bench=. helpers
PASS
BenchmarkChannels 100000 15560 ns/op
BenchmarkMutex 1000000 2669 ns/op
ok helpers 4.452s
Run Code Online (Sandbox Code Playgroud)
我不会期望互斥体会那么快......
进一步的想法
Nic*_*ood 19
不要在链接页面中使用sync/atomic
Package atomic提供了用于实现同步算法的低级原子内存原语.这些功能需要非常小心才能正确使用.除了特殊的低级应用程序,最好使用通道或同步包的功能来实现同步
上次我不得不这样做时,我用一个互斥量来看看你的第二个例子,看起来像你的第三个带有频道的例子.当事情变得非常繁忙时,频道代码赢了,但要确保你的频道缓冲区很大.
Aeq*_*tas 17
如果你正在尝试同步一个工作池(例如允许n goroutines来处理一些工作量),那么通道是一个非常好的方法,但如果你真正需要的只是一个计数器(例如页面浏览量) )然后他们是矫枉过正的.在同步和同步/原子包在那里帮助.
import "sync/atomic"
type count32 int32
func (c *count32) increment() int32 {
return atomic.AddInt32((*int32)(c), 1)
}
func (c *count32) get() int32 {
return atomic.LoadInt32((*int32)(c))
}
Run Code Online (Sandbox Code Playgroud)
不要害怕使用互斥锁和锁只是因为你认为它们"不适合Go".在你的第二个例子中,它绝对清楚发生了什么,这很重要.您将不得不亲自尝试看看互斥体是多么满足,以及添加复杂性是否会提高性能.
如果你确实需要提高性能,也许分片是最好的方法:http: //play.golang.org/p/uLirjskGeN
缺点是您的计数只会与分片决定时一样最新.调用可能还会有很多性能time.Since(),但是,一如既往地先测量它:)
使用同步/原子的另一个答案适用于页面计数器之类的东西,但不适用于向外部 API 提交唯一标识符。为此,您需要一个“增量并返回”操作,该操作只能作为 CAS 循环实现。
这是一个围绕 int32 的 CAS 循环来生成唯一的消息 ID:
import "sync/atomic"
type UniqueID struct {
counter int32
}
func (c *UniqueID) Get() int32 {
for {
val := atomic.LoadInt32(&c.counter)
if atomic.CompareAndSwapInt32(&c.counter, val, val+1) {
return val
}
}
}
Run Code Online (Sandbox Code Playgroud)
要使用它,只需执行以下操作:
requestID := client.msgID.Get()
form.Set("id", requestID)
Run Code Online (Sandbox Code Playgroud)
这比通道有一个优势,因为它不需要那么多额外的空闲资源 - 使用现有的 goroutines 来请求 ID,而不是为程序需要的每个计数器使用一个 goroutine。
TODO:针对渠道进行基准测试。我将猜测通道在无争用情况下更糟,在高争用情况下更好,因为它们需要排队,而此代码只是为了赢得比赛而旋转。
老问题,但我只是偶然发现了这一点,它可能会有所帮助:https : //github.com/uber-go/atomic
基本上 Uber 的工程师已经在sync/atomic包之上构建了一些不错的 util 函数
我还没有在生产中对此进行测试,但代码库非常小,大多数功能的实现都非常标准
绝对优于使用通道或基本互斥锁